Skip to main content

mapping/
file.rs

1// Copyright 2026 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use crate::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
24/// Default readahead size used for streaming reads and decompression (128 KiB).
25pub const READ_AHEAD_SIZE: u64 = 128 * 1024;
26
27/// Calculates the readahead size for a given chunk size, rounding down `suggested_read_ahead_size`
28/// to a multiple of `chunk_size`, or returning `chunk_size` if it is larger than
29/// `suggested_read_ahead_size`.
30pub 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/// Transformation applied to stored data before it can be presented to the reader.
39#[derive(Default)]
40pub enum Transform {
41    /// The file is stored uncompressed and unencrypted.
42    #[default]
43    None,
44    /// The file is stored compressed with chunked metadata.
45    Compressed(CompressionInfo),
46    /// The file is stored encrypted with a cipher.
47    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
90/// A mapped file containing extents and an optional transformation (compression or encryption).
91pub 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    /// Returns the extents mapping logical offsets to device offsets.
103    pub fn extents(&self) -> &Extents {
104        &self.extents
105    }
106
107    /// Returns the uncompressed size of the file in bytes.
108    pub fn uncompressed_size(&self) -> u64 {
109        self.uncompressed_size
110    }
111
112    /// Returns the transform applied to the file if compressed or encrypted.
113    pub fn transform(&self) -> &Transform {
114        &self.transform
115    }
116
117    /// Returns decompression metadata if the file is compressed.
118    pub fn compression_info(&self) -> Option<&CompressionInfo> {
119        match &self.transform {
120            Transform::Compressed(info) => Some(info),
121            _ => None,
122        }
123    }
124
125    /// Returns the cipher if the file is encrypted.
126    pub fn cipher(&self) -> Option<&Arc<dyn Cipher>> {
127        match &self.transform {
128            Transform::Encrypted(cipher) => Some(cipher),
129            _ => None,
130        }
131    }
132
133    /// Streams and decodes the uncompressed range requested by `page_request`, applying readahead.
134    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                    // The range must be out of range. This should be handled when `page_request`
200                    // is dropped.
201                    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                    // Both `buffer` and `dest` are aligned to 64 bytes.
236                    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
284/// Trait for handling deliveries (page request buffers and blob registrations) from [`Files`].
285pub trait DeliveryHandler: Send + Sync + 'static {
286    type Request: PageRequest;
287
288    /// Produces a [`PageRequest`] to receive decompressed or uncompressed page data.
289    fn get_page_request(self: &Arc<Self>, key: u64, range: Range<u64>) -> Self::Request;
290
291    /// Registers a blob's Merkle tree leaf hashes with the downstream verifier.
292    fn register_blob(&self, _key: u64, _merkle_leaves: &[[u8; 32]]) -> Result<(), Error> {
293        Ok(())
294    }
295
296    /// Unregisters a file when it is closed.
297    fn unregister_file(&self, _key: u64) {}
298}
299
300/// A no-op [`DeliveryHandler`] for sessions that do not run a kernel pager or verify blobs
301/// (such as intermediate partition mapping sessions).
302pub 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/// An RAII guard for the background pager thread spawned by [`Files::spawn_pager_thread`].
316///
317/// The thread services page requests for as long as this guard is kept alive, and stops
318/// when this guard is dropped.
319#[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
339/// A thread-safe registry of active [`File`] instances indexed by their Zircon pager port key.
340pub 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    /// Creates a new file registry with the provided block service, delivery handler, and port.
350    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    /// Spawns a background thread to service pager requests for `self`.
361    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    // Runs a synchronous event loop that listens on `self.port` for pager page requests and
373    // dispatches them to `Self::handle_page_request`.
374    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    /// Returns a reference to the block service.
398    pub fn service(&self) -> &Arc<S> {
399        &self.service
400    }
401
402    /// Registers a blob's Merkle tree leaf hashes with the downstream delivery handler.
403    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    // Handles a page request from `PagerThread`.
408    //
409    // If the file is loaded, reads the range into a newly allocated buffer immediately.
410    // If the file is currently loading or unmapped, queues the request to be fulfilled
411    // when loaded.
412    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                // Page requests can arrive before the corresponding `Mappings` command
428                // is processed. We create a loading entry and queue the request so it can be
429                // serviced once the mapping arrives and metadata finishes loading.
430                //
431                // Potential weakness: Unknown keys are unbounded. If an invalid or bogus page
432                // request arrives for a key that is never mapped, this entry will remain in
433                // memory until the session is dropped.
434                entry.insert(FileEntry::Loading(LoadingSlot {
435                    requests: vec![req],
436                    waiters: Vec::new(),
437                }));
438            }
439        }
440    }
441
442    /// Marks `key` as currently loading metadata, preserving any page requests that arrived
443    /// prior to the mapping command.
444    pub fn begin_loading(&self, key: u64) {
445        self.map.lock().entry(key).or_default();
446    }
447
448    /// Inserts a file into the registry under `key` and dispatches any page requests that
449    /// arrived while metadata was loading.
450    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            // `insert` is called inside a block I/O completion callback. Calling
466            // `file.read_range` here would deadlock a single-threaded block completion worker
467            // when a read requires multiple buffer allocations from a full pool: allocating the
468            // next chunk blocks the completion thread waiting for earlier chunks to complete
469            // and free their buffers, which that same thread is responsible for completing.
470            // Instead, hand `reqs` to `PagerThread` so the callback can return immediately.
471            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    /// Returns the loaded [`File`] registered under `key`, if present.
484    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    /// Returns a future that completes with the loaded [`File`] registered under `key`,
493    /// waiting asynchronously if the file is currently loading or hasn't arrived yet.
494    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    /// Removes the file registered under `key`.
521    pub fn remove(&self, key: u64) {
522        self.delivery_handler.unregister_file(key);
523        self.map.lock().remove(&key);
524    }
525
526    /// Returns `true` if `key` is currently in the loading state.
527    #[cfg(test)]
528    pub fn is_loading(&self, key: u64) -> bool {
529        matches!(self.map.lock().get(&key), Some(FileEntry::Loading(_)))
530    }
531
532    /// Returns `true` if `key` is currently in the loaded state.
533    #[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    /// Creates a file registry without a pager for intermediate (e.g. partition) sessions.
541    pub fn new_without_pager(service: Arc<S>) -> Self {
542        Self::new(service, NoopDeliveryHandler, zx::Port::create())
543    }
544}
545
546/// The version where `BlobMetadata` was introduced in Fxfs.
547/// Eventually, we'll need to integrate Fxfs's code for upgrading data structures.
548const BLOB_METADATA_VERSION: u32 = 53;
549
550/// Deserializes versioned `BlobMetadata` from the raw on-disk bytes.
551fn 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
569/// Metadata decoded from on-disk storage describing a blob's uncompressed size, compression info,
570/// and Merkle leaves. Compression info will be None if blob is uncompressed.
571pub struct DecodedBlobMetadata {
572    pub uncompressed_size: u64,
573    pub compression_info: Option<CompressionInfo>,
574    pub merkle_leaves: MerkleLeaves,
575}
576
577/// Reads blob metadata asynchronously from `metadata_extents` using `service`.
578/// Once the metadata is retrieved, deserialized, and parsed, `callback` is invoked with
579/// [`DecodedBlobMetadata`].
580///
581/// On failure or if the operation is aborted, `callback` is dropped without being invoked.
582pub 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
674/// RAII guard that manages the lifecycle of a file transitioning from loading metadata to loaded.
675///
676/// When metadata is being fetched asynchronously from storage, the file entry in [`Files`]
677/// remains in the [`FileEntry::Loading`] state, accumulating incoming page requests in its queue.
678///
679/// - On success: [`LoadingFileGuard::commit`] consumes the guard, stores the fully initialized
680///   [`File`], and dispatches all queued page requests to [`PagerThread`] to be fulfilled.
681/// - On failure or cancellation: If dropped before `commit` is called (e.g. due to storage I/O
682///   error, corrupted metadata, or session teardown), the `Drop` implementation cleans up the
683///   entry by removing `key` from [`Files`]. Dropping the loading slot drops all queued
684///   [`PageRequest`] objects, which fails the pending page requests in the kernel pager.
685struct 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    /// Commits the loaded file to the registry, transferring ownership and dispatching any queued
692    /// page requests to [`PagerThread`] to be fulfilled.
693    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
706/// Processes a raw mapping command (`RawMappingCommand`), decoding extent descriptors,
707/// reading blob metadata from storage, registering Merkle leaves via `files`, and
708/// inserting/removing the file from `files`.
709pub 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        // Force a small block allocation limit (e.g. 4096 bytes) so that read_aligned_range
916        // splits the compressed chunks across multiple consecutive OwnedBuffers!
917        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        // start >= end should be a no-op returning Ok(())
945        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        // Force capping to 4096 bytes per buffer allocation so read_range processes
1019        // 4 separate chunks.
1020        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; // 256 KiB
1055        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        // Request 1 block at offset 4096. Readahead should expand to 0..128 KiB.
1066        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; // 256 KiB
1079        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        // Request 1 block at offset 132 KiB (135168..139264).
1090        // Readahead should expand to 128 KiB..256 KiB (131072..262144).
1091        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        // Request 1 block at offset 132 KiB (135168..139264).
1116        // Readahead window starts at 128 KiB (131072) and would normally extend to 256 KiB
1117        // (262144), but should be capped at page_aligned_size (143360).
1118        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        // The remaining tail bytes within the last page must be zero-filled.
1132        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        // Request 1 block in the second 128 KiB readahead window (e.g. 135168..139264).
1176        // Readahead should expand to 131072..262144 (chunks 4, 5, 6, 7).
1177        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        // Pre-fill destination buffer with 0xFF bytes to verify tail zeroing
1232        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        // Trigger delayed storage read completion.
1368        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        // Corrupt CompressionInfo where compressed_offsets are invalid (descending)
1379        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        // Callback should have been dropped on error without sending a message.
1402        assert!(rx.recv().is_err());
1403    }
1404
1405    #[fuchsia::test]
1406    fn test_read_blob_metadata_corrupt_bincode_drops_callback() {
1407        // Corrupt random bytes that cannot be deserialized as BlobMetadata
1408        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        // Callback should have been dropped on error without sending a message.
1420        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        // File is initially in Loading state.
1480        assert!(files.is_loading(99));
1481
1482        // Complete metadata read (which will fail due to corrupt CompressionInfo).
1483        service.wait_and_trigger_sync();
1484
1485        // Files map should have cleaned up the failed entry on drop.
1486        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        // Prepare serialized BlobMetadata with Merkle leaves on a mock block device.
1495        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        // Use DelayedBlockService to hold back the block read completion for the blob's metadata.
1503        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        // Encode extent descriptors for data (block 0) and metadata (block 1).
1524        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        // Inject the Mappings command and extent payload into the VMO FIFO.
1535        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        // Process the mapping command.
1557        let msg = receiver.peek().unwrap();
1558        process_mapping_command(&msg, &files).unwrap();
1559
1560        // While the block read for metadata is pending completion, the file is in the "loading"
1561        // state and leaves are not yet registered.
1562        assert!(files.is_loading(42));
1563        assert!(registered.lock().unwrap().is_none());
1564
1565        // Trigger completion of the pending block read.
1566        service.wait_and_trigger_sync();
1567
1568        // Once the block read completes and metadata is decoded, verify that:
1569        // - The file transitioned out of "loading" and is committed into `Files`.
1570        // - The `register_blob` callback was invoked with the expected key and Merkle leaves.
1571        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        // Pager request arrives while loading. Buffer should be created and held.
1675        files.handle_page_request(99, 0..4096);
1676        assert!(!dropped.load(Ordering::Relaxed));
1677
1678        // Complete metadata read (which will fail due to corrupt CompressionInfo).
1679        service.wait_and_trigger_sync();
1680
1681        // Dropping guard cleans up LoadingSlot, which drops the queued buffer.
1682        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        // Metadata stored at physical block 1
1696        device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1697            .copy_from_slice(&encoded_metadata);
1698        // Data stored at physical block 0
1699        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        // Encode mapping command with 1 data extent and 1 metadata extent
1721        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            // Wait until metadata block read has been submitted to service
1761            // (file is in Loading state).
1762            while service.pending.lock().is_empty() {
1763                std::thread::sleep(std::time::Duration::from_millis(5));
1764            }
1765
1766            // Trigger page fault from a detached background thread while file is loading metadata.
1767            std::thread::spawn(move || {
1768                let mut b = [0u8; 1];
1769                let _ = vmo_blob_clone.read(&mut b, 0);
1770            });
1771
1772            // Complete metadata block read; this will insert the file and immediately
1773            // drain queued page requests.
1774            service.wait_and_trigger_sync();
1775        });
1776
1777        // Complete data block read requested by pager thread when draining queued request.
1778        service.wait_and_trigger_sync();
1779
1780        // Check if page request made during metadata read was serviced!
1781        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        // Metadata stored at physical block 1
1794        device_data[BLOCK_SIZE as usize..BLOCK_SIZE as usize + encoded_metadata.len()]
1795            .copy_from_slice(&encoded_metadata);
1796        // Data stored at physical block 0
1797        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        // Encode mapping command with 1 data extent and 1 metadata extent
1819        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        // Trigger page fault from a background thread BEFORE the mapping command is processed.
1852        std::thread::spawn(move || {
1853            let mut b = [0u8; 1];
1854            let _ = vmo_blob_clone.read(&mut b, 0);
1855        });
1856
1857        // Wait briefly to ensure page request arrives and is queued as pending mapping in `files`.
1858        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            // Trigger metadata read completion.
1868            service.wait_and_trigger_sync();
1869        });
1870
1871        // Complete data block read requested by pager thread when draining queued request.
1872        service.wait_and_trigger_sync();
1873
1874        // Check if page request made before mapping command was serviced!
1875        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        // 1. Send Mappings command with 0 metadata extents (uncompressed file, loads immediately).
1896        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        // 2. Send CloseBlob command.
1922        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        // Valid metadata with version 53.
1964        let bytes = serialize_metadata(&metadata);
1965        let deserialized = deserialize_blob_metadata(&bytes).unwrap();
1966        assert_eq!(deserialized, metadata);
1967
1968        // Buffer too short (< 4 bytes).
1969        assert!(deserialize_blob_metadata(&[1, 2, 3]).is_err());
1970
1971        // Unsupported version (< 53).
1972        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        // Wait before insert.
1986        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        // Insert the file.
1992        files.insert(42, file_clone);
1993        let loaded_file = wait_task.await;
1994        assert_eq!(loaded_file.uncompressed_size(), BLOCK_SIZE);
1995
1996        // Wait after insert.
1997        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; // 256 KiB
2029        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        // Request 1 block at offset 132 KiB (135168..139264).
2046        // Readahead window expands to 128 KiB..256 KiB (131072..262144), ensuring that
2047        // current_offset starts at a non-zero offset (128 KiB) with correct sector tweak.
2048        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        // Capping to BLOCK_SIZE forces read_aligned_range to issue separate buffer reads for each chunk.
2072        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        // 128 KiB of blob data (> 64 KiB buffer pool) + 1 block (4 KiB) of metadata.
2285        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        // 64 KiB buffer pool (< 128 KiB READ_AHEAD_SIZE).
2295        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        // Start metadata read (txn 0: 4 KiB). Allocates 4 KiB from the 64 KiB pool and queues txn 0
2313        // in `block_service.pending`.
2314        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        // Queue a page request while txn 0 is still pending. Because `File::read_range` expands
2331        // page requests to `READ_AHEAD_SIZE` (128 KiB) and the buffer pool is 64 KiB, servicing
2332        // this request requires two sequential 64 KiB buffer allocations (txn 1 and txn 2).
2333        files.handle_page_request(blob_key, 0..4096);
2334
2335        // Complete txn 0 (metadata) followed by txn 1 and txn 2 (the two 64 KiB data chunks) on a
2336        // single completion thread. Completing txn 0 inserts the file, which must not block the
2337        // completion thread while servicing the queued page request.
2338        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        // Verify that the 128 KiB readahead was split into two 64 KiB block reads (matching the
2350        // 64 KiB pool capacity) and both chunks were committed to the pager.
2351        assert_eq!(rx.commits(), vec![(0, 65536), (65536, 65536)]);
2352        // Verify that the full 128 KiB of blob data was copied into the page request's buffer.
2353        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}