1use crate::reader::{BlockService, read_aligned_range};
6use crate::{
7 ENCRYPTION_KEY_SIZE, Extents, MappingCommand, NullPageRequest, PageRequest, RawMappingCommand,
8};
9use anyhow::{Error, anyhow, bail};
10use blob_metadata::{BlobFormat, BlobMetadata, MerkleLeaves};
11use byteorder::{LittleEndian, ReadBytesExt};
12use delivery_blob::compression::{CompressionAlgorithm, CompressionInfo, StreamingDecompressor};
13use fuchsia_sync::Mutex;
14use futures::channel::oneshot;
15use fxfs_crypto::{Cipher, FxfsCipher, UnwrappedKey};
16use std::borrow::Borrow;
17use std::cmp::min;
18use std::collections::hash_map::{Entry, HashMap};
19use std::ops::{ControlFlow, Range};
20use std::sync::Arc;
21use vmo_fifo::Message;
22use zx::sys::zx_page_request_command_t::ZX_PAGER_VMO_READ;
23
24pub const READ_AHEAD_SIZE: u64 = 128 * 1024;
26
27pub fn read_ahead_size_for_chunk_size(chunk_size: u64, suggested_read_ahead_size: u64) -> u64 {
31 if chunk_size >= suggested_read_ahead_size {
32 chunk_size
33 } else {
34 (suggested_read_ahead_size / chunk_size) * chunk_size
35 }
36}
37
38#[derive(Default)]
40pub enum Transform {
41 #[default]
43 None,
44 Compressed(CompressionInfo),
46 Encrypted(Arc<dyn Cipher>),
48}
49
50impl std::fmt::Debug for Transform {
51 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
52 match self {
53 Self::None => write!(f, "Transform::None"),
54 Self::Compressed(_) => write!(f, "Transform::Compressed(..)"),
55 Self::Encrypted(cipher) => f.debug_tuple("Transform::Encrypted").field(cipher).finish(),
56 }
57 }
58}
59
60impl From<Option<CompressionInfo>> for Transform {
61 fn from(info: Option<CompressionInfo>) -> Self {
62 match info {
63 Some(info) => Self::Compressed(info),
64 None => Self::None,
65 }
66 }
67}
68
69impl From<CompressionInfo> for Transform {
70 fn from(info: CompressionInfo) -> Self {
71 Self::Compressed(info)
72 }
73}
74
75impl From<Arc<dyn Cipher>> for Transform {
76 fn from(cipher: Arc<dyn Cipher>) -> Self {
77 Self::Encrypted(cipher)
78 }
79}
80
81#[derive(Clone)]
82struct FileCompression(Arc<File>);
83
84impl Borrow<CompressionInfo> for FileCompression {
85 fn borrow(&self) -> &CompressionInfo {
86 self.0.compression_info().expect("compression_info missing")
87 }
88}
89
90pub struct File {
92 extents: Extents,
93 uncompressed_size: u64,
94 transform: Transform,
95}
96
97impl File {
98 pub fn new(extents: Extents, uncompressed_size: u64, transform: Transform) -> Self {
99 Self { extents, uncompressed_size, transform }
100 }
101
102 pub fn extents(&self) -> &Extents {
104 &self.extents
105 }
106
107 pub fn uncompressed_size(&self) -> u64 {
109 self.uncompressed_size
110 }
111
112 pub fn transform(&self) -> &Transform {
114 &self.transform
115 }
116
117 pub fn compression_info(&self) -> Option<&CompressionInfo> {
119 match &self.transform {
120 Transform::Compressed(info) => Some(info),
121 _ => None,
122 }
123 }
124
125 pub fn cipher(&self) -> Option<&Arc<dyn Cipher>> {
127 match &self.transform {
128 Transform::Encrypted(cipher) => Some(cipher),
129 _ => None,
130 }
131 }
132
133 pub fn read_range(
135 self: &Arc<Self>,
136 service: &(impl BlockService + ?Sized),
137 mut page_request: impl PageRequest,
138 ) {
139 let page_size = zx::system_get_page_size() as u64;
140 let original_range = page_request.range();
141 if original_range.is_empty() {
142 return;
143 }
144
145 let page_aligned_size = self.uncompressed_size.next_multiple_of(page_size);
146 if original_range.start >= page_aligned_size {
147 return;
148 }
149
150 let read_ahead_size = match &self.transform {
151 Transform::Compressed(info) => {
152 read_ahead_size_for_chunk_size(info.chunk_size(), READ_AHEAD_SIZE)
153 }
154 Transform::None | Transform::Encrypted(_) => READ_AHEAD_SIZE,
155 };
156
157 let read_range = (original_range.start / read_ahead_size) * read_ahead_size
158 ..std::cmp::min(
159 original_range.end.next_multiple_of(read_ahead_size),
160 page_aligned_size,
161 );
162 if page_request.prepare(read_range.clone()).is_err() {
163 return;
164 }
165
166 match &self.transform {
167 Transform::None => {
168 let mut current_offset = read_range.start;
169 let uncompressed_size = self.uncompressed_size;
170
171 read_aligned_range(&self.extents, read_range, service, move |res| {
172 let buffer = match res {
173 Ok(buffer) => buffer,
174 Err(error) => {
175 log::error!(error:?; "Failed to read blocks for mapped file");
176 return ControlFlow::Break(());
177 }
178 };
179 let valid_len =
180 min(buffer.len() as u64, uncompressed_size.saturating_sub(current_offset))
181 as usize;
182 let dest = page_request.mut_ptr_slice().subslice_mut(0..buffer.len());
183 let (mut head, mut tail) = dest.split_at_mut(valid_len);
184 head.copy_from_ptr_slice(buffer.as_ptr_slice().subslice(0..valid_len));
185 tail.fill(0);
186 if page_request.commit(buffer.len()).is_err() {
187 return ControlFlow::Break(());
188 }
189 current_offset += buffer.len() as u64;
190 ControlFlow::Continue(())
191 });
192 }
193 Transform::Compressed(_) => {
194 let Ok((mut decompressor, aligned_range)) = StreamingDecompressor::new(
195 FileCompression(self.clone()),
196 self.uncompressed_size,
197 page_request,
198 ) else {
199 return;
202 };
203
204 read_aligned_range(&self.extents, aligned_range, service, move |res| {
205 let buffer = match res {
206 Ok(buffer) => buffer,
207 Err(error) => {
208 log::error!(
209 error:?;
210 "Failed to read blocks for compressed mapped file"
211 );
212 return ControlFlow::Break(());
213 }
214 };
215 if decompressor.push(buffer.as_ptr_slice()).is_err() {
216 return ControlFlow::Break(());
217 }
218 ControlFlow::Continue(())
219 });
220 }
221 Transform::Encrypted(cipher) => {
222 let mut current_offset = read_range.start;
223 let uncompressed_size = self.uncompressed_size;
224 let cipher = cipher.clone();
225
226 read_aligned_range(&self.extents, read_range, service, move |res| {
227 let buffer = match res {
228 Ok(buffer) => buffer,
229 Err(error) => {
230 log::error!(error:?; "Failed to read blocks for encrypted mapped file");
231 return ControlFlow::Break(());
232 }
233 };
234 let mut dest = page_request.mut_ptr_slice().subslice_mut(0..buffer.len());
235 if let Err(error) = cipher.decrypt_to(
237 0,
238 0,
239 0,
240 current_offset,
241 buffer.as_ptr_slice(),
242 dest.reborrow(),
243 ) {
244 log::error!(error:?; "Failed to decrypt buffer");
245 return ControlFlow::Break(());
246 }
247 let valid_len =
248 min(buffer.len() as u64, uncompressed_size.saturating_sub(current_offset))
249 as usize;
250 dest.subslice_mut(valid_len..buffer.len()).fill(0);
251 if page_request.commit(buffer.len()).is_err() {
252 return ControlFlow::Break(());
253 }
254 current_offset += buffer.len() as u64;
255 ControlFlow::Continue(())
256 });
257 }
258 }
259 }
260}
261
262struct LoadingSlot<R> {
263 requests: Vec<R>,
264 waiters: Vec<oneshot::Sender<Arc<File>>>,
265}
266
267impl<R> Default for LoadingSlot<R> {
268 fn default() -> Self {
269 Self { requests: Vec::new(), waiters: Vec::new() }
270 }
271}
272
273enum FileEntry<R> {
274 Loading(LoadingSlot<R>),
275 Loaded(Arc<File>),
276}
277
278impl<R> Default for FileEntry<R> {
279 fn default() -> Self {
280 Self::Loading(LoadingSlot::default())
281 }
282}
283
284pub trait DeliveryHandler: Send + Sync + 'static {
286 type Request: PageRequest;
287
288 fn get_page_request(self: &Arc<Self>, key: u64, range: Range<u64>) -> Self::Request;
290
291 fn register_blob(&self, _key: u64, _merkle_leaves: &[[u8; 32]]) -> Result<(), Error> {
293 Ok(())
294 }
295
296 fn unregister_file(&self, _key: u64) {}
298}
299
300pub struct NoopDeliveryHandler;
303
304impl DeliveryHandler for NoopDeliveryHandler {
305 type Request = NullPageRequest;
306
307 fn get_page_request(self: &Arc<Self>, _key: u64, _range: Range<u64>) -> NullPageRequest {
308 NullPageRequest
309 }
310}
311
312const PAGER_SHUTDOWN_KEY: u64 = 0;
313const PAGER_WAKE_KEY: u64 = 1;
314
315#[must_use = "Dropping PagerThread immediately stops the background pager thread"]
320pub struct PagerThread<S: ?Sized, D: DeliveryHandler> {
321 files: Arc<Files<S, D>>,
322 thread: Option<std::thread::JoinHandle<()>>,
323}
324
325impl<S: ?Sized, D: DeliveryHandler> Drop for PagerThread<S, D> {
326 fn drop(&mut self) {
327 let packet = zx::Packet::from_user_packet(
328 PAGER_SHUTDOWN_KEY,
329 0,
330 zx::UserPacket::from_u8_array([0; 32]),
331 );
332 let _ = self.files.port.queue(&packet);
333 if let Some(thread) = self.thread.take() {
334 let _ = thread.join();
335 }
336 }
337}
338
339pub struct Files<S: ?Sized, D: DeliveryHandler> {
341 service: Arc<S>,
342 delivery_handler: Arc<D>,
343 port: zx::Port,
344 map: Mutex<HashMap<u64, FileEntry<D::Request>>>,
345 deferred_page_requests: Mutex<Vec<(Arc<File>, D::Request)>>,
346}
347
348impl<S: BlockService + ?Sized, D: DeliveryHandler> Files<S, D> {
349 pub fn new(service: Arc<S>, delivery_handler: D, port: zx::Port) -> Self {
351 Self {
352 service,
353 delivery_handler: Arc::new(delivery_handler),
354 port,
355 map: Mutex::new(HashMap::new()),
356 deferred_page_requests: Mutex::new(Vec::new()),
357 }
358 }
359
360 pub fn spawn_pager_thread(self: &Arc<Self>) -> PagerThread<S, D>
362 where
363 S: 'static,
364 {
365 let files = self.clone();
366 let thread = std::thread::spawn(move || {
367 files.run_pager_loop();
368 });
369 PagerThread { files: self.clone(), thread: Some(thread) }
370 }
371
372 fn run_pager_loop(&self) {
375 while let Ok(packet) = self.port.wait(zx::MonotonicInstant::INFINITE) {
376 match packet.contents() {
377 zx::PacketContents::User(_) => match packet.key() {
378 PAGER_SHUTDOWN_KEY => break,
379 PAGER_WAKE_KEY => {
380 let requests = std::mem::take(&mut *self.deferred_page_requests.lock());
381 for (file, req) in requests {
382 file.read_range(self.service.as_ref(), req);
383 }
384 }
385 _ => {}
386 },
387 zx::PacketContents::Pager(pager_packet)
388 if pager_packet.command() == ZX_PAGER_VMO_READ =>
389 {
390 self.handle_page_request(packet.key(), pager_packet.range());
391 }
392 _ => {}
393 }
394 }
395 }
396
397 pub fn service(&self) -> &Arc<S> {
399 &self.service
400 }
401
402 pub fn register_blob(&self, key: u64, merkle_leaves: &[[u8; 32]]) -> Result<(), Error> {
404 self.delivery_handler.register_blob(key, merkle_leaves)
405 }
406
407 fn handle_page_request(&self, key: u64, range: Range<u64>) {
413 let req = self.delivery_handler.get_page_request(key, range);
414 let mut map = self.map.lock();
415 match map.entry(key) {
416 Entry::Occupied(mut entry) => match entry.get_mut() {
417 FileEntry::Loaded(file) => {
418 let file = Arc::clone(file);
419 drop(map);
420 file.read_range(self.service.as_ref(), req);
421 }
422 FileEntry::Loading(slot) => {
423 slot.requests.push(req);
424 }
425 },
426 Entry::Vacant(entry) => {
427 entry.insert(FileEntry::Loading(LoadingSlot {
435 requests: vec![req],
436 waiters: Vec::new(),
437 }));
438 }
439 }
440 }
441
442 pub fn begin_loading(&self, key: u64) {
445 self.map.lock().entry(key).or_default();
446 }
447
448 fn insert(&self, key: u64, file: Arc<File>) {
451 let (reqs, waiters) = {
452 let mut map = self.map.lock();
453 let prev = map.insert(key, FileEntry::Loaded(file.clone()));
454 match prev {
455 Some(FileEntry::Loading(slot)) => (slot.requests, slot.waiters),
456 _ => (Vec::new(), Vec::new()),
457 }
458 };
459
460 for waiter in waiters {
461 let _ = waiter.send(file.clone());
462 }
463
464 if !reqs.is_empty() {
465 self.deferred_page_requests
472 .lock()
473 .extend(reqs.into_iter().map(|req| (file.clone(), req)));
474 let packet = zx::Packet::from_user_packet(
475 PAGER_WAKE_KEY,
476 0,
477 zx::UserPacket::from_u8_array([0; 32]),
478 );
479 let _ = self.port.queue(&packet);
480 }
481 }
482
483 pub fn get_file(&self, key: u64) -> Option<Arc<File>> {
485 let map = self.map.lock();
486 match map.get(&key) {
487 Some(FileEntry::Loaded(file)) => Some(file.clone()),
488 _ => None,
489 }
490 }
491
492 pub async fn wait_for_file(&self, key: u64) -> Result<Arc<File>, Error> {
495 let receiver = {
496 let mut map = self.map.lock();
497 match map.entry(key) {
498 Entry::Occupied(mut entry) => match entry.get_mut() {
499 FileEntry::Loaded(file) => return Ok(file.clone()),
500 FileEntry::Loading(slot) => {
501 let (sender, receiver) = oneshot::channel();
502 slot.waiters.push(sender);
503 receiver
504 }
505 },
506 Entry::Vacant(entry) => {
507 let (sender, receiver) = oneshot::channel();
508 entry.insert(FileEntry::Loading(LoadingSlot {
509 requests: Vec::new(),
510 waiters: vec![sender],
511 }));
512 receiver
513 }
514 }
515 };
516
517 receiver.await.map_err(|_| anyhow!("File loading cancelled"))
518 }
519
520 pub fn remove(&self, key: u64) {
522 self.delivery_handler.unregister_file(key);
523 self.map.lock().remove(&key);
524 }
525
526 #[cfg(test)]
528 pub fn is_loading(&self, key: u64) -> bool {
529 matches!(self.map.lock().get(&key), Some(FileEntry::Loading(_)))
530 }
531
532 #[cfg(test)]
534 pub fn is_loaded(&self, key: u64) -> bool {
535 matches!(self.map.lock().get(&key), Some(FileEntry::Loaded(_)))
536 }
537}
538
539impl<S: BlockService + ?Sized> Files<S, NoopDeliveryHandler> {
540 pub fn new_without_pager(service: Arc<S>) -> Self {
542 Self::new(service, NoopDeliveryHandler, zx::Port::create())
543 }
544}
545
546const BLOB_METADATA_VERSION: u32 = 53;
549
550fn deserialize_blob_metadata(mut bytes: &[u8]) -> Result<BlobMetadata, anyhow::Error> {
552 use bincode::Options;
553 let options = bincode::DefaultOptions::new().allow_trailing_bytes();
554
555 let version = bytes.read_u32::<LittleEndian>()?;
556 if version < BLOB_METADATA_VERSION {
557 bail!(
558 "Unsupported blob metadata version: {} (expected >= {})",
559 version,
560 BLOB_METADATA_VERSION
561 );
562 }
563
564 options
565 .deserialize::<BlobMetadata>(bytes)
566 .map_err(|e| anyhow!("Failed to deserialize BlobMetadata: {e:?}"))
567}
568
569pub struct DecodedBlobMetadata {
572 pub uncompressed_size: u64,
573 pub compression_info: Option<CompressionInfo>,
574 pub merkle_leaves: MerkleLeaves,
575}
576
577pub fn read_blob_metadata(
583 service: &(impl BlockService + ?Sized),
584 metadata_extents: &Extents,
585 stored_data_size: u64,
586 callback: impl FnOnce(DecodedBlobMetadata) + Send + 'static,
587) {
588 let mut total_metadata_len = 0u64;
589 for extent in metadata_extents.iter_extents(0) {
590 total_metadata_len += extent.len();
591 }
592 if total_metadata_len == 0 {
593 callback(DecodedBlobMetadata {
594 uncompressed_size: stored_data_size,
595 compression_info: None,
596 merkle_leaves: MerkleLeaves::new(),
597 });
598 return;
599 }
600
601 let mut callback = Some(callback);
602 let mut metadata_bytes = Some(Vec::with_capacity(total_metadata_len as usize));
603 read_aligned_range(metadata_extents, 0..total_metadata_len, service, move |res| {
604 let Ok(buffer) = res else {
605 log::error!(error:? = res.unwrap_err(); "Failed to read metadata blocks");
606 return ControlFlow::Break(());
607 };
608 let bytes = metadata_bytes.as_mut().unwrap();
609 let copy_len = min(buffer.len(), total_metadata_len as usize - bytes.len());
610 buffer.as_ptr_slice().subslice(0..copy_len).append_to(bytes);
611 if bytes.len() < total_metadata_len as usize {
612 return ControlFlow::Continue(());
613 }
614
615 let metadata_bytes = metadata_bytes.take().unwrap();
616 let cb = callback.take().unwrap();
617 let metadata = match deserialize_blob_metadata(&metadata_bytes) {
618 Ok(metadata) => metadata,
619 Err(error) => {
620 log::error!(error:?; "Failed to deserialize BlobMetadata");
621 return ControlFlow::Break(());
622 }
623 };
624
625 let merkle_leaves = metadata.merkle_leaves;
626 let res = match metadata.format {
627 BlobFormat::Uncompressed => DecodedBlobMetadata {
628 uncompressed_size: stored_data_size,
629 compression_info: None,
630 merkle_leaves,
631 },
632 BlobFormat::ChunkedZstd { uncompressed_size, chunk_size, compressed_offsets } => {
633 match CompressionInfo::new(
634 chunk_size,
635 stored_data_size,
636 &compressed_offsets,
637 CompressionAlgorithm::Zstd,
638 ) {
639 Ok(compression_info) => DecodedBlobMetadata {
640 uncompressed_size,
641 compression_info: Some(compression_info),
642 merkle_leaves,
643 },
644 Err(error) => {
645 log::error!(error:?; "Failed to parse Zstd CompressionInfo");
646 return ControlFlow::Break(());
647 }
648 }
649 }
650 BlobFormat::ChunkedLz4 { uncompressed_size, chunk_size, compressed_offsets } => {
651 match CompressionInfo::new(
652 chunk_size,
653 stored_data_size,
654 &compressed_offsets,
655 CompressionAlgorithm::Lz4,
656 ) {
657 Ok(compression_info) => DecodedBlobMetadata {
658 uncompressed_size,
659 compression_info: Some(compression_info),
660 merkle_leaves,
661 },
662 Err(error) => {
663 log::error!(error:?; "Failed to parse Lz4 CompressionInfo");
664 return ControlFlow::Break(());
665 }
666 }
667 }
668 };
669 cb(res);
670 ControlFlow::Break(())
671 });
672}
673
674struct LoadingFileGuard<S: BlockService + ?Sized + 'static, D: DeliveryHandler> {
686 files: Option<Arc<Files<S, D>>>,
687 key: u64,
688}
689
690impl<S: BlockService + ?Sized + 'static, D: DeliveryHandler> LoadingFileGuard<S, D> {
691 fn commit(mut self, file: Arc<File>) {
694 self.files.take().unwrap().insert(self.key, file);
695 }
696}
697
698impl<S: BlockService + ?Sized + 'static, D: DeliveryHandler> Drop for LoadingFileGuard<S, D> {
699 fn drop(&mut self) {
700 if let Some(files) = self.files.take() {
701 files.remove(self.key);
702 }
703 }
704}
705
706pub fn process_mapping_command<S: BlockService + ?Sized + 'static, D: DeliveryHandler>(
710 msg: &Message<'_, RawMappingCommand>,
711 files: &Arc<Files<S, D>>,
712) -> Result<(), Error> {
713 let cmd = **msg;
714 match MappingCommand::try_from(cmd)? {
715 MappingCommand::Mappings {
716 key,
717 offset,
718 stored_size,
719 device_offset,
720 metadata_count,
721 extent_count,
722 encrypted,
723 } => {
724 if encrypted && metadata_count > 0 {
725 bail!("Encrypted files cannot have blob metadata");
726 }
727 let extent_bytes_len = (extent_count as usize)
728 .checked_mul(8)
729 .ok_or_else(|| anyhow!("Overflow calculating extent byte length"))?;
730 let metadata_bytes_len = (metadata_count as usize)
731 .checked_mul(8)
732 .ok_or_else(|| anyhow!("Overflow calculating metadata extent byte length"))?;
733 let key_bytes_len = if encrypted { ENCRYPTION_KEY_SIZE } else { 0 };
734 let total_bytes_len = extent_bytes_len
735 .checked_add(metadata_bytes_len)
736 .and_then(|len| len.checked_add(key_bytes_len))
737 .ok_or_else(|| anyhow!("Overflow calculating total extent byte length"))?;
738 let payload_len: u32 = total_bytes_len
739 .try_into()
740 .map_err(|_| anyhow!("Extent byte length exceeds u32"))?;
741
742 let payload_bytes = msg.payload_slice(offset, payload_len);
743 let data_bytes = payload_bytes.subslice(0..extent_bytes_len);
744 let data_extents = Extents::from_encoded(data_bytes.iter_as::<u64>(), device_offset)
745 .ok_or_else(|| anyhow!("Failed to decode data extents"))?;
746
747 if encrypted {
748 let key_bytes = payload_bytes.subslice(extent_bytes_len..total_bytes_len);
749 let cipher_key = UnwrappedKey::new(key_bytes.to_vec());
750 let cipher = Arc::new(FxfsCipher::new(&cipher_key)) as Arc<dyn Cipher>;
751 let file =
752 Arc::new(File::new(data_extents, stored_size, Transform::Encrypted(cipher)));
753 files.insert(key, file);
754 return Ok(());
755 }
756
757 let metadata_bytes = payload_bytes.subslice(extent_bytes_len..total_bytes_len);
758 let metadata_extents =
759 Extents::from_encoded(metadata_bytes.iter_as::<u64>(), device_offset)
760 .ok_or_else(|| anyhow!("Failed to decode metadata extents"))?;
761
762 files.begin_loading(key);
763 let service = files.service().clone();
764 let guard = LoadingFileGuard { files: Some(files.clone()), key };
765 let files_clone = files.clone();
766 read_blob_metadata(service.as_ref(), &metadata_extents, stored_size, move |metadata| {
767 if let Err(error) = files_clone.register_blob(key, &metadata.merkle_leaves) {
768 log::error!(error:?; "Failed to register blob {key}");
769 return;
770 }
771 let file = Arc::new(File::new(
772 data_extents,
773 metadata.uncompressed_size,
774 metadata.compression_info.into(),
775 ));
776 guard.commit(file);
777 });
778 Ok(())
779 }
780 MappingCommand::CloseBlob { key } => {
781 files.remove(key);
782 Ok(())
783 }
784 }
785}
786
787#[cfg(test)]
788mod tests {
789 use super::*;
790 use crate::reader::tests::FakeBlockService;
791 use crate::reader::{MAX_READ_BUFFER_SIZE, OwnedBuffer};
792 use crate::testing::{TestDeliveryHandler, TestVecBuffer};
793 use crate::{BLOCK_SIZE, Extent};
794 use anyhow::Error;
795 use bincode::Options;
796 use byteorder::WriteBytesExt;
797 use delivery_blob::compression::{ChunkedArchiveOptions, CompressionAlgorithm};
798 use fuchsia_async as fasync;
799 use fxfs_crypto::{Cipher, FxfsCipher, UnwrappedKey};
800 use std::sync::Arc;
801 use storage_device::buffer_allocator::{BufferAllocator, BufferSource};
802 use storage_ptr_slice::MutPtrByteSlice;
803
804 fn serialize_metadata(metadata: &BlobMetadata) -> Vec<u8> {
805 let mut bytes = Vec::new();
806 bytes.write_u32::<LittleEndian>(BLOB_METADATA_VERSION).unwrap();
807 bincode::DefaultOptions::new()
808 .allow_trailing_bytes()
809 .serialize_into(&mut bytes, metadata)
810 .unwrap();
811 bytes
812 }
813
814 #[test]
815 fn test_read_range_uncompressed() {
816 let block_count = 8;
817 let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
818 for (i, byte) in expected_data.iter_mut().enumerate() {
819 *byte = (i % 255) as u8;
820 }
821 let service = FakeBlockService::new(expected_data.clone());
822
823 let extents = Extents::try_new([Extent::new(0..(8 * BLOCK_SIZE), Some(0))], 0).unwrap();
824 let file = Arc::new(File::new(extents, 8 * BLOCK_SIZE, Transform::None));
825
826 let (page_request, rx) = TestVecBuffer::new_with_range(0..(8 * BLOCK_SIZE));
827 file.read_range(&service, page_request);
828
829 assert_eq!(rx.commits(), vec![(0, (8 * BLOCK_SIZE) as usize)]);
830 assert_eq!(rx.output(), expected_data);
831 }
832
833 #[test]
834 fn test_read_range_compressed_zstd() {
835 let uncompressed_size = 32768 * 2 + 1024;
836 let mut uncompressed_data = vec![0u8; uncompressed_size];
837 for (i, byte) in uncompressed_data.iter_mut().enumerate() {
838 *byte = ((i * 7) % 255) as u8;
839 }
840
841 let options =
842 ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Zstd };
843 let archive =
844 delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
845
846 let mut compressed_offsets = vec![0];
847 let mut compressed_data = vec![];
848 for chunk in archive.chunks() {
849 compressed_data.extend_from_slice(&chunk.compressed_data);
850 compressed_offsets.push(compressed_data.len() as u64);
851 }
852 compressed_offsets.pop();
853
854 let chunk_size = archive.chunk_size();
855 let stored_size = compressed_data.len() as u64;
856 let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
857 let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
858 device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
859 let service = FakeBlockService::new(device_data);
860
861 let extents =
862 Extents::try_new([Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
863 let compression_info = CompressionInfo::new(
864 chunk_size as u64,
865 stored_size,
866 &compressed_offsets,
867 CompressionAlgorithm::Zstd,
868 )
869 .unwrap();
870 let file = Arc::new(File::new(extents, uncompressed_size as u64, compression_info.into()));
871
872 let dest_alloc_size = uncompressed_size.next_multiple_of(chunk_size);
873 let (mut page_request, rx) = TestVecBuffer::new_with_range(0..(uncompressed_size as u64));
874 page_request.data.resize(dest_alloc_size, 0);
875 file.read_range(&service, page_request);
876
877 assert_eq!(
878 rx.commits(),
879 vec![
880 (0, chunk_size),
881 (chunk_size as u64, chunk_size),
882 (chunk_size as u64 * 2, chunk_size)
883 ]
884 );
885 assert_eq!(&rx.output()[..uncompressed_size], &uncompressed_data[..]);
886 }
887
888 #[test]
889 fn test_read_range_compressed_lz4_split_across_buffers() {
890 let uncompressed_size = 32768 * 2;
891 let mut uncompressed_data = vec![0u8; uncompressed_size];
892 for (i, byte) in uncompressed_data.iter_mut().enumerate() {
893 *byte = ((i * 13) % 255) as u8;
894 }
895
896 let options =
897 ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Lz4 };
898 let archive =
899 delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
900
901 let mut compressed_offsets = vec![0];
902 let mut compressed_data = vec![];
903 for chunk in archive.chunks() {
904 compressed_data.extend_from_slice(&chunk.compressed_data);
905 compressed_offsets.push(compressed_data.len() as u64);
906 }
907 compressed_offsets.pop();
908
909 let chunk_size = archive.chunk_size();
910 let stored_size = compressed_data.len() as u64;
911 let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
912 let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
913 device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
914
915 let service = FakeBlockService::new_with_cap(device_data, Some(4096));
918
919 let extents =
920 Extents::try_new([Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
921 let compression_info = CompressionInfo::new(
922 chunk_size as u64,
923 stored_size,
924 &compressed_offsets,
925 CompressionAlgorithm::Lz4,
926 )
927 .unwrap();
928 let file = Arc::new(File::new(extents, uncompressed_size as u64, compression_info.into()));
929
930 let (page_request, rx) = TestVecBuffer::new_with_range(0..(uncompressed_size as u64));
931 file.read_range(&service, page_request);
932
933 assert_eq!(rx.commits(), vec![(0, chunk_size), (chunk_size as u64, chunk_size)]);
934 assert_eq!(rx.output(), uncompressed_data);
935 }
936
937 #[test]
938 fn test_read_range_invalid_range_noop() {
939 let service = FakeBlockService::new(vec![0u8; 8192]);
940 let extents = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
941 let file = Arc::new(File::new(extents, 8192, Transform::None));
942
943 let (page_request, rx) = TestVecBuffer::new_with_range(4096..4096);
944 file.read_range(&service, page_request);
946 assert_eq!(rx.commits().len(), 0);
947 }
948
949 #[test]
950 fn test_file_getters() {
951 let extents = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
952 let uncompressed_size = 8192u64;
953
954 let file_uncompressed = File::new(extents, uncompressed_size, Transform::None);
955 assert_eq!(file_uncompressed.uncompressed_size(), 8192);
956 assert!(matches!(file_uncompressed.transform(), Transform::None));
957 assert!(file_uncompressed.compression_info().is_none());
958 assert!(file_uncompressed.cipher().is_none());
959
960 let compression_info =
961 CompressionInfo::new(32768, 4096, &[0], CompressionAlgorithm::Zstd).unwrap();
962 let file_compressed = File::new(
963 Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap(),
964 uncompressed_size,
965 compression_info.into(),
966 );
967 assert!(matches!(file_compressed.transform(), Transform::Compressed(_)));
968 assert!(file_compressed.compression_info().is_some());
969 assert!(file_compressed.cipher().is_none());
970
971 let key = UnwrappedKey::new(vec![0x42u8; 32]);
972 let cipher: Arc<dyn Cipher> = Arc::new(FxfsCipher::new(&key));
973 let file_encrypted = File::new(
974 Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap(),
975 uncompressed_size,
976 cipher.into(),
977 );
978 assert!(matches!(file_encrypted.transform(), Transform::Encrypted(_)));
979 assert!(file_encrypted.compression_info().is_none());
980 assert!(file_encrypted.cipher().is_some());
981 }
982
983 #[test]
984 fn test_read_range_block_service_error_returns_err() {
985 struct FailingBlockService;
986 impl BlockService for FailingBlockService {
987 fn allocate_buffer(&self, max_len: usize) -> storage_device::buffer::OwnedBuffer {
988 FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
989 }
990 fn read_blocks(
991 &self,
992 _device_offset: u64,
993 _dest_buffer: storage_device::buffer::OwnedBuffer,
994 _on_complete: Box<
995 dyn FnOnce(Result<storage_device::buffer::OwnedBuffer, Error>) + Send,
996 >,
997 ) -> Result<(), Error> {
998 Err(anyhow::anyhow!("block read failure"))
999 }
1000 }
1001
1002 let extents = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
1003 let file = Arc::new(File::new(extents, 8192, Transform::None));
1004
1005 let (page_request, rx) = TestVecBuffer::new_with_range(0..8192);
1006
1007 file.read_range(&FailingBlockService, page_request);
1008 assert_eq!(rx.commits().len(), 0);
1009 }
1010
1011 #[test]
1012 fn test_read_range_uncompressed_multi_chunk() {
1013 let block_count = 4;
1014 let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
1015 for (i, byte) in expected_data.iter_mut().enumerate() {
1016 *byte = ((i * 11) % 255) as u8;
1017 }
1018 let service = FakeBlockService::new_with_cap(expected_data.clone(), Some(4096));
1021
1022 let extents =
1023 Extents::try_new([Extent::new(0..(block_count * BLOCK_SIZE), Some(0))], 0).unwrap();
1024 let file = Arc::new(File::new(extents, block_count * BLOCK_SIZE, Transform::None));
1025
1026 let (page_request, rx) = TestVecBuffer::new_with_range(0..(block_count * BLOCK_SIZE));
1027 file.read_range(&service, page_request);
1028
1029 assert_eq!(rx.commits().len(), 4);
1030 assert_eq!(rx.output(), expected_data);
1031 }
1032
1033 #[test]
1034 fn test_read_range_uncompressed_unaligned_uncompressed_size() {
1035 let uncompressed_size = 5000u64;
1036 let mut expected_data = vec![0u8; 8192];
1037 for (i, byte) in expected_data.iter_mut().enumerate() {
1038 *byte = (i % 251) as u8;
1039 }
1040 let service = FakeBlockService::new(expected_data.clone());
1041
1042 let extents = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
1043 let file = Arc::new(File::new(extents, uncompressed_size, Transform::None));
1044
1045 let (page_request, rx) = TestVecBuffer::new_with_range(0..8192);
1046 file.read_range(&service, page_request);
1047
1048 assert_eq!(rx.commits(), vec![(0, 8192)]);
1049 assert_eq!(&rx.output()[..5000], &expected_data[..5000]);
1050 }
1051
1052 #[test]
1053 fn test_read_range_uncompressed_readahead() {
1054 let total_blocks = 64; let mut expected_data = vec![0u8; (total_blocks * BLOCK_SIZE) as usize];
1056 for (i, byte) in expected_data.iter_mut().enumerate() {
1057 *byte = ((i * 17) % 251) as u8;
1058 }
1059 let service = FakeBlockService::new(expected_data.clone());
1060
1061 let extents =
1062 Extents::try_new([Extent::new(0..(total_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
1063 let file = Arc::new(File::new(extents, total_blocks * BLOCK_SIZE, Transform::None));
1064
1065 let (page_request, rx) = TestVecBuffer::new_with_range(4096..8192);
1067 file.read_range(&service, page_request);
1068
1069 assert_eq!(rx.commits(), vec![(0, READ_AHEAD_SIZE as usize)]);
1070 assert_eq!(
1071 &rx.output()[..READ_AHEAD_SIZE as usize],
1072 &expected_data[..READ_AHEAD_SIZE as usize]
1073 );
1074 }
1075
1076 #[test]
1077 fn test_read_range_uncompressed_readahead_second_window() {
1078 let total_blocks = 64; let mut expected_data = vec![0u8; (total_blocks * BLOCK_SIZE) as usize];
1080 for (i, byte) in expected_data.iter_mut().enumerate() {
1081 *byte = ((i * 19) % 251) as u8;
1082 }
1083 let service = FakeBlockService::new(expected_data.clone());
1084
1085 let extents =
1086 Extents::try_new([Extent::new(0..(total_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
1087 let file = Arc::new(File::new(extents, total_blocks * BLOCK_SIZE, Transform::None));
1088
1089 let (page_request, rx) = TestVecBuffer::new_with_range(135168..139264);
1092 file.read_range(&service, page_request);
1093
1094 assert_eq!(rx.commits(), vec![(READ_AHEAD_SIZE, READ_AHEAD_SIZE as usize)]);
1095 assert_eq!(
1096 &rx.output()[..READ_AHEAD_SIZE as usize],
1097 &expected_data[READ_AHEAD_SIZE as usize..2 * READ_AHEAD_SIZE as usize]
1098 );
1099 }
1100
1101 #[test]
1102 fn test_read_range_uncompressed_readahead_tail_capped() {
1103 let uncompressed_size = 140_000u64;
1104 let total_blocks = (uncompressed_size.next_multiple_of(BLOCK_SIZE) / BLOCK_SIZE) as u64;
1105 let mut expected_data = vec![0u8; (total_blocks * BLOCK_SIZE) as usize];
1106 for i in 0..uncompressed_size as usize {
1107 expected_data[i] = ((i * 23) % 251) as u8;
1108 }
1109 let service = FakeBlockService::new(expected_data.clone());
1110
1111 let extents =
1112 Extents::try_new([Extent::new(0..(total_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
1113 let file = Arc::new(File::new(extents, uncompressed_size, Transform::None));
1114
1115 let (page_request, rx) = TestVecBuffer::new_with_range(135168..139264);
1119 file.read_range(&service, page_request);
1120
1121 let expected_start = READ_AHEAD_SIZE;
1122 let page_aligned_size = uncompressed_size.next_multiple_of(BLOCK_SIZE);
1123 let expected_len = (page_aligned_size - expected_start) as usize;
1124 assert_eq!(rx.commits(), vec![(expected_start, expected_len)]);
1125
1126 let valid_len = (uncompressed_size - expected_start) as usize;
1127 assert_eq!(
1128 &rx.output()[..valid_len],
1129 &expected_data[expected_start as usize..uncompressed_size as usize]
1130 );
1131 assert!(rx.output()[valid_len..expected_len].iter().all(|&b| b == 0));
1133 }
1134
1135 #[test]
1136 fn test_read_range_compressed_second_readahead_window() {
1137 let chunk_count = 8;
1138 let chunk_size = 32768usize;
1139 let uncompressed_size = chunk_count * chunk_size;
1140 let mut uncompressed_data = vec![0u8; uncompressed_size];
1141 for (i, byte) in uncompressed_data.iter_mut().enumerate() {
1142 *byte = ((i * 7) % 251) as u8;
1143 }
1144
1145 let options =
1146 ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Zstd };
1147 let archive =
1148 delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
1149
1150 let mut compressed_offsets = vec![0];
1151 let mut compressed_data = vec![];
1152 for chunk in archive.chunks() {
1153 compressed_data.extend_from_slice(&chunk.compressed_data);
1154 compressed_offsets.push(compressed_data.len() as u64);
1155 }
1156 compressed_offsets.pop();
1157
1158 let stored_size = compressed_data.len() as u64;
1159 let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
1160 let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
1161 device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
1162 let service = FakeBlockService::new(device_data);
1163
1164 let extents =
1165 Extents::try_new([Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
1166 let compression_info = CompressionInfo::new(
1167 chunk_size as u64,
1168 stored_size,
1169 &compressed_offsets,
1170 CompressionAlgorithm::Zstd,
1171 )
1172 .unwrap();
1173 let file = Arc::new(File::new(extents, uncompressed_size as u64, compression_info.into()));
1174
1175 let (page_request, rx) = TestVecBuffer::new_with_range(135168..139264);
1178 file.read_range(&service, page_request);
1179
1180 assert_eq!(
1181 rx.commits(),
1182 vec![
1183 (131072, chunk_size),
1184 (163840, chunk_size),
1185 (196608, chunk_size),
1186 (229376, chunk_size),
1187 ]
1188 );
1189 assert_eq!(&rx.output()[..131072], &uncompressed_data[131072..262144]);
1190 }
1191
1192 #[test]
1193 fn test_read_range_compressed_partial_final_chunk_zero_tail() {
1194 let uncompressed_size = 32768 + 1024;
1195 let mut uncompressed_data = vec![0u8; uncompressed_size];
1196 for (i, byte) in uncompressed_data.iter_mut().enumerate() {
1197 *byte = ((i * 13) % 251) as u8;
1198 }
1199
1200 let options =
1201 ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Zstd };
1202 let archive =
1203 delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
1204
1205 let mut compressed_offsets = vec![0];
1206 let mut compressed_data = vec![];
1207 for chunk in archive.chunks() {
1208 compressed_data.extend_from_slice(&chunk.compressed_data);
1209 compressed_offsets.push(compressed_data.len() as u64);
1210 }
1211 compressed_offsets.pop();
1212
1213 let chunk_size = archive.chunk_size();
1214 let stored_size = compressed_data.len() as u64;
1215 let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
1216 let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
1217 device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
1218 let service = FakeBlockService::new(device_data);
1219
1220 let extents =
1221 Extents::try_new([Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
1222 let compression_info = CompressionInfo::new(
1223 chunk_size as u64,
1224 stored_size,
1225 &compressed_offsets,
1226 CompressionAlgorithm::Zstd,
1227 )
1228 .unwrap();
1229 let file = Arc::new(File::new(extents, uncompressed_size as u64, compression_info.into()));
1230
1231 let (mut page_request, rx) = TestVecBuffer::new_with_range(0..(uncompressed_size as u64));
1233 page_request.data.resize(65536, 0);
1234 page_request.data.fill(0xFF);
1235 file.read_range(&service, page_request);
1236
1237 assert_eq!(rx.commits(), vec![(0, chunk_size), (chunk_size as u64, chunk_size)]);
1238 assert_eq!(&rx.output()[..uncompressed_size], &uncompressed_data[..]);
1239 assert_eq!(&rx.output()[uncompressed_size..65536], &[0u8; 31744]);
1240 }
1241
1242 #[test]
1243 fn test_files_registry() {
1244 let extents = Extents::try_new([Extent::new(0..4096, Some(0))], 0).unwrap();
1245 let file = Arc::new(File::new(extents, 4096, Transform::None));
1246 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
1247 let files = Files::new(
1248 service,
1249 TestDeliveryHandler(|_key, _range| TestVecBuffer::new(4096).0),
1250 zx::Port::create(),
1251 );
1252
1253 assert!(!files.is_loading(100));
1254 files.begin_loading(100);
1255 assert!(files.is_loading(100));
1256
1257 files.insert(100, file.clone());
1258 assert!(!files.is_loading(100));
1259 files.remove(100);
1260 assert!(!files.is_loading(100));
1261 }
1262
1263 #[test]
1264 fn test_files_new_without_pager() {
1265 let extents = Extents::try_new([Extent::new(0..4096, Some(0))], 0).unwrap();
1266 let file = Arc::new(File::new(extents, 4096, Transform::None));
1267 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
1268 let files = Files::new_without_pager(service);
1269
1270 files.insert(100, file.clone());
1271 assert_eq!(files.get_file(100).unwrap().uncompressed_size(), 4096);
1272 }
1273
1274 struct DelayedBlockService {
1275 device_data: Vec<u8>,
1276 pending: Mutex<Vec<Box<dyn FnOnce() + Send>>>,
1277 allocator: Option<Arc<BufferAllocator>>,
1278 }
1279
1280 impl DelayedBlockService {
1281 fn new(device_data: Vec<u8>) -> Arc<Self> {
1282 Arc::new(Self { device_data, pending: Mutex::new(Vec::new()), allocator: None })
1283 }
1284
1285 fn new_with_pool_capacity(device_data: Vec<u8>, pool_capacity: usize) -> Arc<Self> {
1286 Arc::new(Self {
1287 device_data,
1288 pending: Mutex::new(Vec::new()),
1289 allocator: Some(Arc::new(BufferAllocator::new(
1290 BLOCK_SIZE as usize,
1291 BufferSource::new(pool_capacity),
1292 ))),
1293 })
1294 }
1295
1296 fn wait_and_trigger_sync(&self) {
1297 loop {
1298 let callbacks = std::mem::take(&mut *self.pending.lock());
1299 if !callbacks.is_empty() {
1300 for cb in callbacks {
1301 cb();
1302 }
1303 break;
1304 }
1305 std::thread::sleep(std::time::Duration::from_millis(10));
1306 }
1307 }
1308 }
1309
1310 impl BlockService for DelayedBlockService {
1311 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1312 if let Some(allocator) = &self.allocator {
1313 let max_len = std::cmp::min(
1314 std::cmp::min(max_len, MAX_READ_BUFFER_SIZE),
1315 allocator.buffer_source().size(),
1316 );
1317 allocator.allocate_buffer_sync_owned(max_len)
1318 } else {
1319 FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
1320 }
1321 }
1322
1323 fn read_blocks(
1324 &self,
1325 device_offset: u64,
1326 mut dest_buffer: OwnedBuffer,
1327 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1328 ) -> Result<(), Error> {
1329 let start = device_offset as usize;
1330 let end = (start + dest_buffer.len()).min(self.device_data.len());
1331 dest_buffer
1332 .as_mut_ptr_slice()
1333 .subslice_mut(0..end - start)
1334 .copy_from_slice(&self.device_data[start..end]);
1335 self.pending.lock().push(Box::new(move || {
1336 on_complete(Ok(dest_buffer));
1337 }));
1338 Ok(())
1339 }
1340 }
1341
1342 #[fuchsia::test]
1343 fn test_read_blob_metadata_callback() {
1344 let uncompressed_size = 65536u64;
1345 let chunk_size = 32768u64;
1346 let compressed_offsets = vec![0u64, 1200u64];
1347 let metadata = BlobMetadata {
1348 merkle_leaves: MerkleLeaves::new(),
1349 format: BlobFormat::ChunkedZstd {
1350 uncompressed_size,
1351 chunk_size,
1352 compressed_offsets: compressed_offsets.clone(),
1353 },
1354 };
1355 let encoded_metadata = serialize_metadata(&metadata);
1356 let mut device_data = vec![0u8; BLOCK_SIZE as usize];
1357 device_data[..encoded_metadata.len()].copy_from_slice(&encoded_metadata);
1358
1359 let service = DelayedBlockService::new(device_data);
1360 let metadata_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1361
1362 let (tx, rx) = std::sync::mpsc::channel();
1363 read_blob_metadata(service.as_ref(), &metadata_extents, 2400, move |res| {
1364 let _ = tx.send(res);
1365 });
1366
1367 service.wait_and_trigger_sync();
1369
1370 let metadata = rx.recv().unwrap();
1371 assert_eq!(metadata.uncompressed_size, uncompressed_size);
1372 assert!(metadata.compression_info.is_some());
1373 assert_eq!(metadata.merkle_leaves, MerkleLeaves::new());
1374 }
1375
1376 #[fuchsia::test]
1377 fn test_read_blob_metadata_error_drops_callback() {
1378 let metadata = BlobMetadata {
1380 merkle_leaves: MerkleLeaves::new(),
1381 format: BlobFormat::ChunkedZstd {
1382 uncompressed_size: 65536,
1383 chunk_size: 32768,
1384 compressed_offsets: vec![1000, 500],
1385 },
1386 };
1387 let encoded_metadata = serialize_metadata(&metadata);
1388 let mut device_data = vec![0u8; BLOCK_SIZE as usize];
1389 device_data[..encoded_metadata.len()].copy_from_slice(&encoded_metadata);
1390
1391 let service = DelayedBlockService::new(device_data);
1392 let metadata_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1393
1394 let (tx, rx) = std::sync::mpsc::channel();
1395 read_blob_metadata(service.as_ref(), &metadata_extents, 1200, move |res| {
1396 let _ = tx.send(res);
1397 });
1398
1399 service.wait_and_trigger_sync();
1400
1401 assert!(rx.recv().is_err());
1403 }
1404
1405 #[fuchsia::test]
1406 fn test_read_blob_metadata_corrupt_bincode_drops_callback() {
1407 let device_data = vec![0xFFu8; BLOCK_SIZE as usize];
1409 let service = DelayedBlockService::new(device_data);
1410 let metadata_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1411
1412 let (tx, rx) = std::sync::mpsc::channel();
1413 read_blob_metadata(service.as_ref(), &metadata_extents, 1200, move |res| {
1414 let _ = tx.send(res);
1415 });
1416
1417 service.wait_and_trigger_sync();
1418
1419 assert!(rx.recv().is_err());
1421 }
1422
1423 #[fuchsia::test]
1424 fn test_process_mapping_command_error_cleans_up_blob() {
1425 let metadata = BlobMetadata {
1426 merkle_leaves: MerkleLeaves::new(),
1427 format: BlobFormat::ChunkedZstd {
1428 uncompressed_size: 65536,
1429 chunk_size: 32768,
1430 compressed_offsets: vec![1000, 500],
1431 },
1432 };
1433 let encoded_metadata = serialize_metadata(&metadata);
1434 let mut device_data = vec![0u8; (2 * BLOCK_SIZE) as usize];
1435 device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1436 .copy_from_slice(&encoded_metadata);
1437
1438 let service = DelayedBlockService::new(device_data);
1439 let files = Arc::new(Files::new(
1440 service.clone(),
1441 TestDeliveryHandler(|_k, _r| TestVecBuffer::new(4096).0),
1442 zx::Port::create(),
1443 ));
1444
1445 let data_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1446 let meta_extents =
1447 Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(BLOCK_SIZE))], 0).unwrap();
1448 let mut payload_bytes = Vec::new();
1449 for w in
1450 Extents::encode_extents(&data_extents).chain(Extents::encode_extents(&meta_extents))
1451 {
1452 payload_bytes.extend_from_slice(&w.to_le_bytes());
1453 }
1454
1455 let vmo = zx::Vmo::create(65536).unwrap();
1456 let mut sender = vmo_fifo::SyncSender::<crate::RawMappingCommand>::new(
1457 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
1458 1024,
1459 16,
1460 )
1461 .unwrap();
1462 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
1463 payload_buf.data().copy_from_slice(&payload_bytes);
1464 let cmd = crate::RawMappingCommand {
1465 opcode: crate::MAPPINGS_COMMAND,
1466 offset: payload_buf.offset(),
1467 key: 99,
1468 stored_size: 4096,
1469 device_offset: 0,
1470 metadata_count: 1,
1471 extent_count: 1,
1472 };
1473 payload_buf.commit(cmd).unwrap();
1474 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
1475
1476 let msg = receiver.peek().unwrap();
1477 process_mapping_command(&msg, &files).unwrap();
1478
1479 assert!(files.is_loading(99));
1481
1482 service.wait_and_trigger_sync();
1484
1485 assert!(!files.is_loading(99));
1487 }
1488
1489 #[fuchsia::test]
1490 fn test_process_mapping_command_registers_blob_merkle_leaves() {
1491 let leaf1 = [1u8; 32];
1492 let leaf2 = [2u8; 32];
1493
1494 let metadata =
1496 BlobMetadata { merkle_leaves: vec![leaf1, leaf2], format: BlobFormat::Uncompressed };
1497 let encoded_metadata = serialize_metadata(&metadata);
1498 let mut device_data = vec![0u8; (2 * BLOCK_SIZE) as usize];
1499 device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1500 .copy_from_slice(&encoded_metadata);
1501
1502 let service = DelayedBlockService::new(device_data);
1504 let registered = Arc::new(std::sync::Mutex::new(None));
1505 let registered_clone = registered.clone();
1506 struct TestRegisterHandler(Arc<std::sync::Mutex<Option<(u64, Vec<[u8; 32]>)>>>);
1507 impl DeliveryHandler for TestRegisterHandler {
1508 type Request = TestVecBuffer;
1509 fn get_page_request(self: &Arc<Self>, _key: u64, _range: Range<u64>) -> Self::Request {
1510 TestVecBuffer::new(4096).0
1511 }
1512 fn register_blob(&self, key: u64, leaves: &[[u8; 32]]) -> Result<(), Error> {
1513 *self.0.lock().unwrap() = Some((key, leaves.to_vec()));
1514 Ok(())
1515 }
1516 }
1517 let files = Arc::new(Files::new(
1518 service.clone(),
1519 TestRegisterHandler(registered_clone),
1520 zx::Port::create(),
1521 ));
1522
1523 let data_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1525 let meta_extents =
1526 Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(BLOCK_SIZE))], 0).unwrap();
1527 let data_extent_words = Extents::encode_extents(&data_extents);
1528 let meta_extent_words = Extents::encode_extents(&meta_extents);
1529 let mut payload_bytes = Vec::new();
1530 for w in data_extent_words.chain(meta_extent_words) {
1531 payload_bytes.extend_from_slice(&w.to_le_bytes());
1532 }
1533
1534 let vmo = zx::Vmo::create(65536).unwrap();
1536 let mut sender = vmo_fifo::SyncSender::<crate::RawMappingCommand>::new(
1537 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
1538 1024,
1539 16,
1540 )
1541 .unwrap();
1542 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
1543 payload_buf.data().copy_from_slice(&payload_bytes);
1544 let cmd = crate::RawMappingCommand {
1545 opcode: crate::MAPPINGS_COMMAND,
1546 offset: payload_buf.offset(),
1547 key: 42,
1548 stored_size: 4096,
1549 device_offset: 0,
1550 metadata_count: 1,
1551 extent_count: 1,
1552 };
1553 payload_buf.commit(cmd).unwrap();
1554 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
1555
1556 let msg = receiver.peek().unwrap();
1558 process_mapping_command(&msg, &files).unwrap();
1559
1560 assert!(files.is_loading(42));
1563 assert!(registered.lock().unwrap().is_none());
1564
1565 service.wait_and_trigger_sync();
1567
1568 assert!(!files.is_loading(42));
1572 assert!(files.get_file(42).is_some());
1573 let (reg_key, reg_leaves) =
1574 registered.lock().unwrap().take().expect("registered callback called");
1575 assert_eq!(reg_key, 42);
1576 assert_eq!(reg_leaves, vec![leaf1, leaf2]);
1577 }
1578
1579 #[fuchsia::test]
1580 fn test_loading_page_request_buffer_dropped_on_failure() {
1581 use delivery_blob::compression::{ChunkedArchiveError, DataBuffer};
1582 use std::sync::atomic::{AtomicBool, Ordering};
1583 use storage_ptr_slice::MutPtrByteSlice;
1584
1585 struct DroppingBuffer {
1586 data: Vec<u8>,
1587 range: Range<u64>,
1588 dropped: Arc<AtomicBool>,
1589 }
1590 impl Drop for DroppingBuffer {
1591 fn drop(&mut self) {
1592 self.dropped.store(true, Ordering::Relaxed);
1593 }
1594 }
1595 impl DataBuffer for DroppingBuffer {
1596 fn range(&self) -> Range<u64> {
1597 self.range.clone()
1598 }
1599 fn mut_ptr_slice(&mut self) -> MutPtrByteSlice<'_> {
1600 MutPtrByteSlice::from(&mut self.data[..])
1601 }
1602 fn commit(&mut self, _size: usize) -> Result<(), ChunkedArchiveError> {
1603 Ok(())
1604 }
1605 }
1606 impl PageRequest for DroppingBuffer {
1607 fn prepare(&mut self, _range: Range<u64>) -> Result<(), ChunkedArchiveError> {
1608 Ok(())
1609 }
1610 }
1611
1612 let metadata = BlobMetadata {
1613 merkle_leaves: MerkleLeaves::new(),
1614 format: BlobFormat::ChunkedZstd {
1615 uncompressed_size: 65536,
1616 chunk_size: 32768,
1617 compressed_offsets: vec![1000, 500],
1618 },
1619 };
1620 let encoded_metadata = serialize_metadata(&metadata);
1621 let mut device_data = vec![0u8; (2 * BLOCK_SIZE) as usize];
1622 device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1623 .copy_from_slice(&encoded_metadata);
1624
1625 let service = DelayedBlockService::new(device_data);
1626 let dropped = Arc::new(AtomicBool::new(false));
1627 let dropped_clone = dropped.clone();
1628 let files = Arc::new(Files::new(
1629 service.clone(),
1630 TestDeliveryHandler(move |_k, r| DroppingBuffer {
1631 data: vec![0u8; 4096],
1632 range: r,
1633 dropped: dropped_clone.clone(),
1634 }),
1635 zx::Port::create(),
1636 ));
1637
1638 let data_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1639 let meta_extents =
1640 Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(BLOCK_SIZE))], 0).unwrap();
1641 let mut payload_bytes = Vec::new();
1642 for w in
1643 Extents::encode_extents(&data_extents).chain(Extents::encode_extents(&meta_extents))
1644 {
1645 payload_bytes.extend_from_slice(&w.to_le_bytes());
1646 }
1647
1648 let vmo = zx::Vmo::create(65536).unwrap();
1649 let mut sender = vmo_fifo::SyncSender::<crate::RawMappingCommand>::new(
1650 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
1651 1024,
1652 16,
1653 )
1654 .unwrap();
1655 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
1656 payload_buf.data().copy_from_slice(&payload_bytes);
1657 let cmd = crate::RawMappingCommand {
1658 opcode: crate::MAPPINGS_COMMAND,
1659 offset: payload_buf.offset(),
1660 key: 99,
1661 stored_size: 4096,
1662 device_offset: 0,
1663 metadata_count: 1,
1664 extent_count: 1,
1665 };
1666 payload_buf.commit(cmd).unwrap();
1667 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
1668
1669 let msg = receiver.peek().unwrap();
1670 process_mapping_command(&msg, &files).unwrap();
1671
1672 assert!(files.is_loading(99));
1673
1674 files.handle_page_request(99, 0..4096);
1676 assert!(!dropped.load(Ordering::Relaxed));
1677
1678 service.wait_and_trigger_sync();
1680
1681 assert!(dropped.load(Ordering::Relaxed));
1683 assert!(!files.is_loading(99));
1684 }
1685
1686 #[fuchsia::test]
1687 async fn test_pager_concurrent_request_during_metadata_read() {
1688 use vmo_fifo::SyncSender;
1689 use zx::{Pager, PagerOptions, Port, Rights, Vmo, VmoOptions};
1690
1691 let metadata =
1692 BlobMetadata { merkle_leaves: MerkleLeaves::new(), format: BlobFormat::Uncompressed };
1693 let encoded_metadata = serialize_metadata(&metadata);
1694 let mut device_data = vec![0u8; (2 * BLOCK_SIZE) as usize];
1695 device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1697 .copy_from_slice(&encoded_metadata);
1698 device_data[..4].copy_from_slice(&[10, 20, 30, 40]);
1700
1701 let service = DelayedBlockService::new(device_data);
1702 let port = Port::create();
1703
1704 let (page_request, rx) = TestVecBuffer::new(BLOCK_SIZE as usize);
1705 let page_request_holder = Arc::new(Mutex::new(Some(page_request)));
1706 let page_request_clone = page_request_holder.clone();
1707
1708 let files = Arc::new(Files::new(
1709 service.clone(),
1710 TestDeliveryHandler(move |_key, _range| page_request_clone.lock().take().unwrap()),
1711 port.duplicate_handle(Rights::SAME_RIGHTS).unwrap(),
1712 ));
1713
1714 let _pager_thread = files.spawn_pager_thread();
1715
1716 let pager = Pager::create(PagerOptions::empty()).unwrap();
1717 let vmo_blob = pager.create_vmo(VmoOptions::empty(), &port, 42, BLOCK_SIZE).unwrap();
1718 let vmo_blob_clone = vmo_blob.duplicate_handle(Rights::SAME_RIGHTS).unwrap();
1719
1720 let data_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1722 let meta_extents =
1723 Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(BLOCK_SIZE))], 0).unwrap();
1724 let mut payload_bytes = Vec::new();
1725 for w in
1726 Extents::encode_extents(&data_extents).chain(Extents::encode_extents(&meta_extents))
1727 {
1728 payload_bytes.extend_from_slice(&w.to_le_bytes());
1729 }
1730
1731 let cmd = crate::RawMappingCommand {
1732 opcode: crate::MAPPINGS_COMMAND,
1733 offset: 0,
1734 key: 42,
1735 stored_size: BLOCK_SIZE,
1736 device_offset: 0,
1737 metadata_count: 1,
1738 extent_count: 1,
1739 };
1740
1741 let vmo = Vmo::create(65536).unwrap();
1742 let mut sender = SyncSender::<crate::RawMappingCommand>::new(
1743 vmo.duplicate_handle(Rights::SAME_RIGHTS).unwrap(),
1744 1024,
1745 16,
1746 )
1747 .unwrap();
1748 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
1749 payload_buf.data().copy_from_slice(&payload_bytes);
1750 payload_buf.commit(cmd).unwrap();
1751 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
1752
1753 std::thread::scope(|s| {
1754 let msg = receiver.peek().unwrap();
1755 let files_for_process = files.clone();
1756 s.spawn(move || {
1757 process_mapping_command(&msg, &files_for_process).unwrap();
1758 });
1759
1760 while service.pending.lock().is_empty() {
1763 std::thread::sleep(std::time::Duration::from_millis(5));
1764 }
1765
1766 std::thread::spawn(move || {
1768 let mut b = [0u8; 1];
1769 let _ = vmo_blob_clone.read(&mut b, 0);
1770 });
1771
1772 service.wait_and_trigger_sync();
1775 });
1776
1777 service.wait_and_trigger_sync();
1779
1780 assert!(!rx.commits().is_empty(), "Page request made during metadata read was dropped!");
1782 }
1783
1784 #[fuchsia::test]
1785 async fn test_pager_request_arrives_before_mapping_command() {
1786 use vmo_fifo::SyncSender;
1787 use zx::{Pager, PagerOptions, Port, Rights, Vmo, VmoOptions};
1788
1789 let metadata =
1790 BlobMetadata { merkle_leaves: MerkleLeaves::new(), format: BlobFormat::Uncompressed };
1791 let encoded_metadata = serialize_metadata(&metadata);
1792 let mut device_data = vec![0u8; (2 * BLOCK_SIZE) as usize];
1793 device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1795 .copy_from_slice(&encoded_metadata);
1796 device_data[..4].copy_from_slice(&[10, 20, 30, 40]);
1798
1799 let service = DelayedBlockService::new(device_data);
1800 let port = Port::create();
1801
1802 let (page_request, rx) = TestVecBuffer::new(BLOCK_SIZE as usize);
1803 let page_request_holder = Arc::new(Mutex::new(Some(page_request)));
1804 let page_request_clone = page_request_holder.clone();
1805
1806 let files = Arc::new(Files::new(
1807 service.clone(),
1808 TestDeliveryHandler(move |_key, _range| page_request_clone.lock().take().unwrap()),
1809 port.duplicate_handle(Rights::SAME_RIGHTS).unwrap(),
1810 ));
1811
1812 let _pager_thread = files.spawn_pager_thread();
1813
1814 let pager = Pager::create(PagerOptions::empty()).unwrap();
1815 let vmo_blob = pager.create_vmo(VmoOptions::empty(), &port, 42, BLOCK_SIZE).unwrap();
1816 let vmo_blob_clone = vmo_blob.duplicate_handle(Rights::SAME_RIGHTS).unwrap();
1817
1818 let data_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1820 let meta_extents =
1821 Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(BLOCK_SIZE))], 0).unwrap();
1822 let mut payload_bytes = Vec::new();
1823 for w in
1824 Extents::encode_extents(&data_extents).chain(Extents::encode_extents(&meta_extents))
1825 {
1826 payload_bytes.extend_from_slice(&w.to_le_bytes());
1827 }
1828
1829 let cmd = crate::RawMappingCommand {
1830 opcode: crate::MAPPINGS_COMMAND,
1831 offset: 0,
1832 key: 42,
1833 stored_size: BLOCK_SIZE,
1834 device_offset: 0,
1835 metadata_count: 1,
1836 extent_count: 1,
1837 };
1838
1839 let vmo = Vmo::create(65536).unwrap();
1840 let mut sender = SyncSender::<crate::RawMappingCommand>::new(
1841 vmo.duplicate_handle(Rights::SAME_RIGHTS).unwrap(),
1842 1024,
1843 16,
1844 )
1845 .unwrap();
1846 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
1847 payload_buf.data().copy_from_slice(&payload_bytes);
1848 payload_buf.commit(cmd).unwrap();
1849 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
1850
1851 std::thread::spawn(move || {
1853 let mut b = [0u8; 1];
1854 let _ = vmo_blob_clone.read(&mut b, 0);
1855 });
1856
1857 std::thread::sleep(std::time::Duration::from_millis(50));
1859
1860 std::thread::scope(|s| {
1861 let msg = receiver.peek().unwrap();
1862 let files_for_process = files.clone();
1863 s.spawn(move || {
1864 process_mapping_command(&msg, &files_for_process).unwrap();
1865 });
1866
1867 service.wait_and_trigger_sync();
1869 });
1870
1871 service.wait_and_trigger_sync();
1873
1874 assert!(!rx.commits().is_empty(), "Page request made before mapping command was dropped!");
1876 }
1877
1878 #[fuchsia::test]
1879 fn test_process_mapping_command_close_blob() {
1880 let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
1881 let files = Arc::new(Files::new(
1882 service,
1883 TestDeliveryHandler(|_k, _r| TestVecBuffer::new(4096).0),
1884 zx::Port::create(),
1885 ));
1886
1887 let vmo = zx::Vmo::create(65536).unwrap();
1888 let mut sender = vmo_fifo::SyncSender::<crate::RawMappingCommand>::new(
1889 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
1890 1024,
1891 16,
1892 )
1893 .unwrap();
1894
1895 let data_extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1897 let mut payload_bytes = Vec::new();
1898 for w in Extents::encode_extents(&data_extents) {
1899 payload_bytes.extend_from_slice(&w.to_le_bytes());
1900 }
1901 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
1902 payload_buf.data().copy_from_slice(&payload_bytes);
1903 let cmd = crate::RawMappingCommand {
1904 opcode: crate::MAPPINGS_COMMAND,
1905 offset: payload_buf.offset(),
1906 key: 123,
1907 stored_size: 4096,
1908 device_offset: 0,
1909 metadata_count: 0,
1910 extent_count: 1,
1911 };
1912 payload_buf.commit(cmd).unwrap();
1913
1914 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
1915 let msg = receiver.peek().unwrap();
1916 process_mapping_command(&msg, &files).unwrap();
1917 msg.pop().unwrap();
1918
1919 assert!(files.is_loaded(123));
1920
1921 sender
1923 .push(crate::RawMappingCommand {
1924 opcode: crate::CLOSE_BLOB_COMMAND,
1925 offset: 0,
1926 key: 123,
1927 stored_size: 0,
1928 device_offset: 0,
1929 metadata_count: 0,
1930 extent_count: 0,
1931 })
1932 .unwrap();
1933
1934 let msg = receiver.peek().unwrap();
1935 process_mapping_command(&msg, &files).unwrap();
1936 msg.pop().unwrap();
1937
1938 assert!(!files.is_loaded(123));
1939 }
1940
1941 #[test]
1942 fn test_read_ahead_size_for_chunk_size() {
1943 assert_eq!(read_ahead_size_for_chunk_size(32 * 1024, 32 * 1024), 32 * 1024);
1944 assert_eq!(read_ahead_size_for_chunk_size(48 * 1024, 32 * 1024), 48 * 1024);
1945 assert_eq!(read_ahead_size_for_chunk_size(64 * 1024, 32 * 1024), 64 * 1024);
1946
1947 assert_eq!(read_ahead_size_for_chunk_size(32 * 1024, 64 * 1024), 64 * 1024);
1948 assert_eq!(read_ahead_size_for_chunk_size(48 * 1024, 64 * 1024), 48 * 1024);
1949 assert_eq!(read_ahead_size_for_chunk_size(64 * 1024, 64 * 1024), 64 * 1024);
1950 assert_eq!(read_ahead_size_for_chunk_size(96 * 1024, 64 * 1024), 96 * 1024);
1951
1952 assert_eq!(read_ahead_size_for_chunk_size(32 * 1024, 128 * 1024), 128 * 1024);
1953 assert_eq!(read_ahead_size_for_chunk_size(48 * 1024, 128 * 1024), 96 * 1024);
1954 assert_eq!(read_ahead_size_for_chunk_size(64 * 1024, 128 * 1024), 128 * 1024);
1955 assert_eq!(read_ahead_size_for_chunk_size(96 * 1024, 128 * 1024), 96 * 1024);
1956 }
1957
1958 #[test]
1959 fn test_deserialize_blob_metadata() {
1960 let metadata =
1961 BlobMetadata { merkle_leaves: MerkleLeaves::new(), format: BlobFormat::Uncompressed };
1962
1963 let bytes = serialize_metadata(&metadata);
1965 let deserialized = deserialize_blob_metadata(&bytes).unwrap();
1966 assert_eq!(deserialized, metadata);
1967
1968 assert!(deserialize_blob_metadata(&[1, 2, 3]).is_err());
1970
1971 let mut old_version_bytes = bytes.clone();
1973 (&mut old_version_bytes[..4]).write_u32::<LittleEndian>(52).unwrap();
1974 assert!(deserialize_blob_metadata(&old_version_bytes).is_err());
1975 }
1976
1977 #[fuchsia::test]
1978 async fn test_wait_for_file_before_and_after_insert() {
1979 let service = Arc::new(FakeBlockService::new(vec![]));
1980 let files = Arc::new(Files::new_without_pager(service));
1981
1982 let extents = Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(0))], 0).unwrap();
1983 let file = Arc::new(File::new(extents, BLOCK_SIZE, Transform::None));
1984
1985 let files_clone = files.clone();
1987 let file_clone = file.clone();
1988 let wait_task =
1989 fasync::Task::spawn(async move { files_clone.wait_for_file(42).await.unwrap() });
1990
1991 files.insert(42, file_clone);
1993 let loaded_file = wait_task.await;
1994 assert_eq!(loaded_file.uncompressed_size(), BLOCK_SIZE);
1995
1996 let loaded_file2 = files.wait_for_file(42).await.unwrap();
1998 assert_eq!(loaded_file2.uncompressed_size(), BLOCK_SIZE);
1999 }
2000
2001 #[test]
2002 fn test_read_range_encrypted() {
2003 let block_count = 8;
2004 let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
2005 for (i, byte) in expected_data.iter_mut().enumerate() {
2006 *byte = (i % 255) as u8;
2007 }
2008
2009 let key = UnwrappedKey::new(vec![0x42u8; 32]);
2010 let cipher: Arc<dyn Cipher> = Arc::new(FxfsCipher::new(&key));
2011
2012 let mut encrypted_data = expected_data.clone();
2013 cipher.encrypt(0, 0, 0, 0, MutPtrByteSlice::from(&mut encrypted_data[..])).unwrap();
2014 let service = FakeBlockService::new(encrypted_data);
2015
2016 let extents = Extents::try_new([Extent::new(0..(8 * BLOCK_SIZE), Some(0))], 0).unwrap();
2017 let file = Arc::new(File::new(extents, 8 * BLOCK_SIZE, cipher.into()));
2018
2019 let (page_request, rx) = TestVecBuffer::new_with_range(0..(8 * BLOCK_SIZE));
2020 file.read_range(&service, page_request);
2021
2022 assert_eq!(rx.commits(), vec![(0, (8 * BLOCK_SIZE) as usize)]);
2023 assert_eq!(rx.output(), expected_data);
2024 }
2025
2026 #[test]
2027 fn test_read_range_encrypted_with_offset() {
2028 let total_blocks = 64; let mut expected_data = vec![0u8; (total_blocks * BLOCK_SIZE) as usize];
2030 for (i, byte) in expected_data.iter_mut().enumerate() {
2031 *byte = ((i * 17) % 255) as u8;
2032 }
2033
2034 let key = UnwrappedKey::new(vec![0x5au8; 32]);
2035 let cipher: Arc<dyn Cipher> = Arc::new(FxfsCipher::new(&key));
2036
2037 let mut encrypted_data = expected_data.clone();
2038 cipher.encrypt(0, 0, 0, 0, MutPtrByteSlice::from(&mut encrypted_data[..])).unwrap();
2039 let service = FakeBlockService::new(encrypted_data);
2040
2041 let extents =
2042 Extents::try_new([Extent::new(0..(total_blocks * BLOCK_SIZE), Some(0))], 0).unwrap();
2043 let file = Arc::new(File::new(extents, total_blocks * BLOCK_SIZE, cipher.into()));
2044
2045 let (page_request, rx) = TestVecBuffer::new_with_range(135168..139264);
2049 file.read_range(&service, page_request);
2050
2051 assert_eq!(rx.commits(), vec![(READ_AHEAD_SIZE, READ_AHEAD_SIZE as usize)]);
2052 assert_eq!(
2053 &rx.output()[..READ_AHEAD_SIZE as usize],
2054 &expected_data[READ_AHEAD_SIZE as usize..2 * READ_AHEAD_SIZE as usize]
2055 );
2056 }
2057
2058 #[test]
2059 fn test_read_range_encrypted_multi_chunk() {
2060 let block_count = 4;
2061 let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
2062 for (i, byte) in expected_data.iter_mut().enumerate() {
2063 *byte = ((i * 31) % 255) as u8;
2064 }
2065
2066 let key = UnwrappedKey::new(vec![0x7fu8; 32]);
2067 let cipher: Arc<dyn Cipher> = Arc::new(FxfsCipher::new(&key));
2068
2069 let mut encrypted_data = expected_data.clone();
2070 cipher.encrypt(0, 0, 0, 0, MutPtrByteSlice::from(&mut encrypted_data[..])).unwrap();
2071 let service = FakeBlockService::new_with_cap(encrypted_data, Some(BLOCK_SIZE as usize));
2073
2074 let extents =
2075 Extents::try_new([Extent::new(0..(block_count * BLOCK_SIZE), Some(0))], 0).unwrap();
2076 let file = Arc::new(File::new(extents, block_count * BLOCK_SIZE, cipher.into()));
2077
2078 let (page_request, rx) = TestVecBuffer::new_with_range(0..(block_count * BLOCK_SIZE));
2079 file.read_range(&service, page_request);
2080
2081 assert_eq!(rx.commits().len(), 4);
2082 assert_eq!(rx.output(), expected_data);
2083 }
2084
2085 #[test]
2086 fn test_read_range_decryption_error() {
2087 #[derive(Debug)]
2088 struct FailingCipher(FxfsCipher);
2089 impl Cipher for FailingCipher {
2090 fn encrypt(
2091 &self,
2092 ino: u64,
2093 attr: u64,
2094 dev: u64,
2095 file: u64,
2096 buf: MutPtrByteSlice<'_>,
2097 ) -> Result<(), Error> {
2098 self.0.encrypt(ino, attr, dev, file, buf)
2099 }
2100 fn decrypt(
2101 &self,
2102 _ino: u64,
2103 _attribute_id: u64,
2104 _device_offset: u64,
2105 _file_offset: u64,
2106 _buffer: MutPtrByteSlice<'_>,
2107 ) -> Result<(), Error> {
2108 bail!("decrypt failed")
2109 }
2110 fn encrypt_filename(&self, ino: u64, name: &mut Vec<u8>) -> Result<(), Error> {
2111 self.0.encrypt_filename(ino, name)
2112 }
2113 fn decrypt_filename(&self, ino: u64, name: &mut Vec<u8>) -> Result<(), Error> {
2114 self.0.decrypt_filename(ino, name)
2115 }
2116 fn hash_code(&self, name: &[u8], filename: &str) -> Option<u32> {
2117 self.0.hash_code(name, filename)
2118 }
2119 fn hash_code_casefold(&self, filename: &str) -> u32 {
2120 self.0.hash_code_casefold(filename)
2121 }
2122 fn supports_inline_encryption(&self) -> bool {
2123 self.0.supports_inline_encryption()
2124 }
2125 fn crypt_ctx(&self, ino: u64, attr: u64, offset: u64) -> Option<(u64, u8)> {
2126 self.0.crypt_ctx(ino, attr, offset)
2127 }
2128 }
2129
2130 let key = UnwrappedKey::new(vec![0x42u8; 32]);
2131 let failing_cipher: Arc<dyn Cipher> = Arc::new(FailingCipher(FxfsCipher::new(&key)));
2132 let service = FakeBlockService::new(vec![0u8; 8192]);
2133 let extents = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
2134 let file = Arc::new(File::new(extents, 8192, failing_cipher.into()));
2135
2136 let (page_request, rx) = TestVecBuffer::new_with_range(0..8192);
2137 file.read_range(&service, page_request);
2138 assert_eq!(rx.commits().len(), 0);
2139 }
2140
2141 #[fuchsia::test]
2142 fn test_process_mapping_command_unencrypted_no_metadata() {
2143 let block_count = 2;
2144 let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
2145 for (i, byte) in expected_data.iter_mut().enumerate() {
2146 *byte = (i % 251) as u8;
2147 }
2148 let service = Arc::new(FakeBlockService::new(expected_data.clone()));
2149 let files = Arc::new(Files::new(
2150 service.clone(),
2151 TestDeliveryHandler(|_k, r| TestVecBuffer::new_with_range(r).0),
2152 zx::Port::create(),
2153 ));
2154
2155 let data_extents =
2156 Extents::try_new([Extent::new(0..(block_count as u64 * BLOCK_SIZE), Some(0))], 0)
2157 .unwrap();
2158 let mut payload_bytes = Vec::new();
2159 for w in Extents::encode_extents(&data_extents) {
2160 payload_bytes.extend_from_slice(&w.to_le_bytes());
2161 }
2162
2163 let vmo = zx::Vmo::create(65536).unwrap();
2164 let mut sender = vmo_fifo::SyncSender::<crate::RawMappingCommand>::new(
2165 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
2166 1024,
2167 16,
2168 )
2169 .unwrap();
2170 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
2171 payload_buf.data().copy_from_slice(&payload_bytes);
2172 let cmd = crate::RawMappingCommand {
2173 opcode: crate::MAPPINGS_COMMAND,
2174 offset: payload_buf.offset(),
2175 key: 100,
2176 stored_size: block_count as u64 * BLOCK_SIZE,
2177 device_offset: 0,
2178 metadata_count: 0,
2179 extent_count: 1,
2180 };
2181 payload_buf.commit(cmd).unwrap();
2182 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
2183
2184 let msg = receiver.peek().unwrap();
2185 process_mapping_command(&msg, &files).unwrap();
2186
2187 let file = files.get_file(100).expect("File should be loaded immediately");
2188 let (page_request, rx) =
2189 TestVecBuffer::new_with_range(0..(block_count as u64 * BLOCK_SIZE));
2190 file.read_range(service.as_ref(), page_request);
2191 assert_eq!(rx.output(), expected_data);
2192 }
2193
2194 #[fuchsia::test]
2195 fn test_process_mapping_command_encrypted() {
2196 struct NoRegisterBlobHandler;
2197 impl DeliveryHandler for NoRegisterBlobHandler {
2198 type Request = TestVecBuffer;
2199 fn get_page_request(self: &Arc<Self>, _key: u64, range: Range<u64>) -> Self::Request {
2200 TestVecBuffer::new_with_range(range).0
2201 }
2202 fn register_blob(&self, _key: u64, _merkle_leaves: &[[u8; 32]]) -> Result<(), Error> {
2203 panic!("register_blob should not be called for encrypted files");
2204 }
2205 }
2206
2207 let block_count = 2;
2208 let mut plaintext = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
2209 for (i, byte) in plaintext.iter_mut().enumerate() {
2210 *byte = (i % 251) as u8;
2211 }
2212
2213 let raw_key = [0x5au8; 32];
2214 let key = UnwrappedKey::new(raw_key.to_vec());
2215 let cipher = FxfsCipher::new(&key);
2216
2217 let mut ciphertext = plaintext.clone();
2218 cipher.encrypt(0, 0, 0, 0, MutPtrByteSlice::from(&mut ciphertext[..])).unwrap();
2219
2220 let service = Arc::new(FakeBlockService::new(ciphertext));
2221 let files =
2222 Arc::new(Files::new(service.clone(), NoRegisterBlobHandler, zx::Port::create()));
2223
2224 let data_extents =
2225 Extents::try_new([Extent::new(0..(block_count as u64 * BLOCK_SIZE), Some(0))], 0)
2226 .unwrap();
2227 let mut payload_bytes = Vec::new();
2228 for w in Extents::encode_extents(&data_extents) {
2229 payload_bytes.extend_from_slice(&w.to_le_bytes());
2230 }
2231 payload_bytes.extend_from_slice(&raw_key);
2232
2233 let vmo = zx::Vmo::create(65536).unwrap();
2234 let mut sender = vmo_fifo::SyncSender::<crate::RawMappingCommand>::new(
2235 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
2236 1024,
2237 16,
2238 )
2239 .unwrap();
2240 let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
2241 payload_buf.data().copy_from_slice(&payload_bytes);
2242 let cmd = crate::RawMappingCommand {
2243 opcode: crate::MAPPINGS_COMMAND | crate::MAPPINGS_FLAG_ENCRYPTED,
2244 offset: payload_buf.offset(),
2245 key: 101,
2246 stored_size: block_count as u64 * BLOCK_SIZE,
2247 device_offset: 0,
2248 metadata_count: 0,
2249 extent_count: 1,
2250 };
2251 payload_buf.commit(cmd).unwrap();
2252
2253 let invalid_cmd = crate::RawMappingCommand {
2254 opcode: crate::MAPPINGS_COMMAND | crate::MAPPINGS_FLAG_ENCRYPTED,
2255 offset: 0,
2256 key: 102,
2257 stored_size: block_count as u64 * BLOCK_SIZE,
2258 device_offset: 0,
2259 metadata_count: 1,
2260 extent_count: 1,
2261 };
2262 sender.reserve_payload(0).unwrap().commit(invalid_cmd).unwrap();
2263
2264 let mut receiver = vmo_fifo::Receiver::<crate::RawMappingCommand>::new(vmo, 16).unwrap();
2265
2266 let msg = receiver.peek().unwrap();
2267 process_mapping_command(&msg, &files).unwrap();
2268 msg.pop().unwrap();
2269
2270 let invalid_msg = receiver.peek().unwrap();
2271 assert!(process_mapping_command(&invalid_msg, &files).is_err());
2272
2273 let file = files.get_file(101).expect("File should be loaded immediately");
2274 assert!(matches!(file.transform(), Transform::Encrypted(_)));
2275
2276 let (page_request, rx) =
2277 TestVecBuffer::new_with_range(0..(block_count as u64 * BLOCK_SIZE));
2278 file.read_range(service.as_ref(), page_request);
2279 assert_eq!(rx.output(), plaintext);
2280 }
2281
2282 #[fuchsia::test]
2283 fn test_queued_page_request_does_not_deadlock_single_completion_thread() {
2284 let data_size = READ_AHEAD_SIZE as usize;
2286 let encoded_metadata = serialize_metadata(&BlobMetadata {
2287 merkle_leaves: MerkleLeaves::new(),
2288 format: BlobFormat::Uncompressed,
2289 });
2290 let mut device_data = vec![0xabu8; data_size + BLOCK_SIZE as usize];
2291 device_data[data_size..data_size + encoded_metadata.len()]
2292 .copy_from_slice(&encoded_metadata);
2293
2294 let block_service = DelayedBlockService::new_with_pool_capacity(device_data, 64 * 1024);
2296
2297 let (page_request, rx) = TestVecBuffer::new_with_range(0..4096);
2298 let page_request = Mutex::new(Some(page_request));
2299 let files = Arc::new(Files::new(
2300 block_service.clone(),
2301 TestDeliveryHandler(move |_key, _range| page_request.lock().take().unwrap()),
2302 zx::Port::create(),
2303 ));
2304 let _pager_thread = files.spawn_pager_thread();
2305
2306 let data_extents = Extents::try_new([Extent::new(0..READ_AHEAD_SIZE, Some(0))], 0).unwrap();
2307 let meta_extents =
2308 Extents::try_new([Extent::new(0..BLOCK_SIZE, Some(READ_AHEAD_SIZE))], 0).unwrap();
2309
2310 let blob_key = 200;
2311
2312 files.begin_loading(blob_key);
2315 let files_clone = files.clone();
2316 read_blob_metadata(
2317 block_service.as_ref(),
2318 &meta_extents,
2319 READ_AHEAD_SIZE,
2320 move |metadata| {
2321 let file = Arc::new(File::new(
2322 data_extents,
2323 metadata.uncompressed_size,
2324 metadata.compression_info.into(),
2325 ));
2326 files_clone.insert(blob_key, file);
2327 },
2328 );
2329
2330 files.handle_page_request(blob_key, 0..4096);
2334
2335 let block_completion_thread = std::thread::spawn({
2339 let block_service = block_service.clone();
2340 move || {
2341 for _ in 0..3 {
2342 block_service.wait_and_trigger_sync();
2343 }
2344 }
2345 });
2346
2347 block_completion_thread.join().unwrap();
2348
2349 assert_eq!(rx.commits(), vec![(0, 65536), (65536, 65536)]);
2352 assert_eq!(rx.output(), vec![0xabu8; data_size]);
2354 }
2355
2356 #[fuchsia::test]
2357 fn test_pager_thread_lifecycle() {
2358 let port = zx::Port::create();
2359 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
2360 let files = Arc::new(Files::new(
2361 service,
2362 TestDeliveryHandler(|_key, _range| TestVecBuffer::new(4096).0),
2363 port,
2364 ));
2365 let thread = files.spawn_pager_thread();
2366 drop(thread);
2367 }
2368
2369 #[fuchsia::test]
2370 fn test_pager_packet() {
2371 let port = zx::Port::create();
2372 let pager = zx::Pager::create(zx::PagerOptions::empty()).expect("create pager");
2373 let vmo = pager.create_vmo(zx::VmoOptions::empty(), &port, 1234, 4096).expect("create vmo");
2374
2375 let vmo_clone = vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate vmo");
2376 let _reader_thread = std::thread::spawn(move || {
2377 let mut b = [0u8; 1];
2378 let _ = vmo_clone.read(&mut b, 0);
2379 });
2380
2381 let packet = port.wait(zx::MonotonicInstant::INFINITE).expect("wait packet");
2382 assert_eq!(packet.key(), 1234);
2383 if let zx::PacketContents::Pager(pager_packet) = packet.contents() {
2384 assert_eq!(pager_packet.command(), ZX_PAGER_VMO_READ);
2385 assert_eq!(pager_packet.range(), 0..4096);
2386 } else {
2387 panic!("Expected pager packet");
2388 }
2389 }
2390}