Skip to main content

blob_pager_and_verifier/
lib.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
5//! Coordinates pager-backed blob VMOs, driver-level page fault handling, and cryptographic
6//! verification.
7//!
8//! # Overview
9//! `BlobPagerAndVerifier` manages pager-backed VMOs for blobs. It decouples page fault servicing
10//! from the filesystem by having the block driver read and decompress raw disk blocks directly,
11//! while `BlobPagerAndVerifier` cryptographically verifies the pages before supplying them to
12//! the pager VMO to unblock client reads.
13//!
14//! # Blob Opening and Caching
15//! When a client requests a blob via `create_vmo`:
16//! - **If already cached**: Clones and returns a new read-only child VMO immediately.
17//! - **If currently opening**: Suspends until the in-flight request finishes, then clones and
18//!   returns a child VMO once ready.
19//! - **If not in cache**: Registers the blob's extent mappings with Fxfs, creates the pager-backed
20//!   VMO, caches it, and returns a child VMO to the caller.
21//!
22//! # Page Fault Handling and Verification
23//! 1. When a client reads an unpopulated page, the kernel generates a page request on the pager
24//!    port, which the block driver monitors directly.
25//! 2. The driver reads and decompresses the disk blocks according to the blob's extent mappings,
26//!    then writes unverified pages and Merkle metadata into a shared delivery queue VMO.
27//! 3. `BlobPagerAndVerifier` verifies the incoming pages against the blob's Merkle tree:
28//!    - On success, it supplies the verified pages to the pager VMO (`zx_pager_supply_pages`),
29//!      unblocking the client read.
30//!    - On verification failure, it fails the request with `ZX_ERR_IO_DATA_INTEGRITY`.
31//!
32//! # Eviction and Teardown
33//! 1. Clients receive read-only child VMOs cloned from the parent pager VMO.
34//! 2. When all client handles to a blob are dropped, the kernel notifies the cache via
35//!    `ZX_VMO_ZERO_CHILDREN`.
36//! 3. The cache evicts the inactive blob and closes its mapping session with Fxfs, reclaiming
37//!    driver extent tracking and kernel pager resources.
38
39mod delivery;
40
41pub use delivery::{
42    DELIVERY_DATA_SIZE, DeliveryQueueProcessor, DeliveryQueueProvider, TestVmoProvider,
43    UnverifiedPages,
44};
45
46use anyhow::{Context, Error, anyhow, bail};
47use event_listener as _;
48use fidl_fuchsia_storage_block as fblock;
49use fidl_fuchsia_storage_mapping as fmapping;
50use fuchsia_async as fasync;
51use fuchsia_hash::{HASH_SIZE, Hash};
52use fuchsia_merkle::{MerkleVerifier, ReadSizedMerkleVerifier};
53use fuchsia_sync::Mutex;
54use std::collections::{HashMap, hash_map};
55use std::ops::Range;
56use std::sync::{Arc, OnceLock, Weak};
57use storage_ptr_slice::PtrByteSlice;
58use zx;
59
60// A `fuchsia_async::PacketReceiver` that watches for a VMO to reach zero children and triggers
61// cache eviction.
62struct ZeroChildrenReceiver {
63    blob: OnceLock<Weak<CachedBlob>>,
64}
65
66impl fasync::PacketReceiver for ZeroChildrenReceiver {
67    fn receive_packet(&self, packet: zx::Packet) {
68        if let zx::PacketContents::SignalOne(signals) = packet.contents() {
69            if signals.observed().contains(zx::Signals::VMO_ZERO_CHILDREN) {
70                if let Some(cached_blob) = self.blob.get().and_then(|w| w.upgrade()) {
71                    if let Some(cache) = cached_blob.cache.upgrade() {
72                        cache.on_zero_children(&cached_blob);
73                    }
74                }
75            }
76        }
77    }
78}
79
80// Holds the state and resources for an open blob residing in the cache.
81//
82// When a blob is opened, this bundles:
83// - The root pager VMO that clients clone from (so we can detect when all handles close).
84// - The Fxfs session key used by the block driver to deliver data and to close the mapping.
85// - The Merkle tree verifier used to check incoming pages from the driver.
86// - The ZeroChildrenReceiver registration that watches for zero children to trigger eviction.
87struct CachedBlob {
88    // The parent pager-backed VMO. Clients only receive child clones so the kernel can notify us
89    // via `ZX_VMO_ZERO_CHILDREN` when all clients have dropped their handles.
90    vmo: zx::Vmo,
91    // The Merkle root hash of the blob (its unique identifier).
92    root_hash: [u8; 32],
93    // Weak back-reference to Cache to avoid an Arc reference cycle, since the cache already
94    // holds an Arc to this blob.
95    cache: Weak<Cache>,
96    // Session key assigned by Fxfs for this blob mapping, used in driver delivery commands and when
97    // closing the session.
98    vmo_key: u64,
99    // Blob size in bytes.
100    len: u64,
101    // Keeps ZeroChildrenReceiver registered on the async loop while this blob is alive.
102    registration: fasync::ReceiverRegistration<ZeroChildrenReceiver>,
103    // Merkle tree verifier, initialised once the driver delivers the blob's leaf hashes.
104    merkle_verifier: OnceLock<ReadSizedMerkleVerifier>,
105}
106
107impl CachedBlob {
108    fn create_child(&self) -> Result<zx::Vmo, zx::Status> {
109        self.vmo.create_child(zx::VmoChildOptions::REFERENCE | zx::VmoChildOptions::NO_WRITE, 0, 0)
110    }
111
112    fn wait_for_zero_children(&self) {
113        let _ = self.vmo.wait_async(
114            fasync::EHandle::local().port(),
115            self.registration.key(),
116            zx::Signals::VMO_ZERO_CHILDREN,
117            zx::WaitAsyncOpts::empty(),
118        );
119    }
120}
121
122impl Drop for CachedBlob {
123    fn drop(&mut self) {
124        if let Some(cache) = self.cache.upgrade() {
125            let session = cache.mapping_session.clone();
126            let vmo_key = self.vmo_key;
127            cache.scope.spawn(async move {
128                let _ = session.close(vmo_key).await;
129            });
130        }
131    }
132}
133
134enum BlobState {
135    // The blob is actively opening: extents are being registered with Fxfs and the pager VMO is
136    // being created. Concurrent requests wait on `event`. Early Merkle leaves delivered by the
137    // driver are placed directly into `verifier`.
138    Opening {
139        root_hash: [u8; 32],
140        event: Arc<event_listener::Event>,
141        verifier: Option<ReadSizedMerkleVerifier>,
142    },
143    // The blob's pager-backed VMO is created and cached. Eviction is triggered when all client
144    // child VMOs are closed (`ZX_VMO_ZERO_CHILDREN`).
145    Ready(Arc<CachedBlob>),
146}
147
148struct CacheState {
149    next_key: u64,
150    // Maps blob Merkle root hash to the client-allocated mapping session key.
151    keys_by_hash: HashMap<[u8; 32], u64>,
152    blob_states_by_key: HashMap<u64, BlobState>,
153}
154
155impl Default for CacheState {
156    fn default() -> Self {
157        Self {
158            next_key: 1,
159            keys_by_hash: HashMap::default(),
160            blob_states_by_key: HashMap::default(),
161        }
162    }
163}
164
165impl CacheState {
166    fn generate_key(&mut self) -> u64 {
167        let key = self.next_key;
168        self.next_key += 1;
169        key
170    }
171}
172
173// Holds a pending reservation for a blob actively being opened.
174//
175// If dropped before `commit` is called (e.g. due to an early FIDL failure, pager VMO creation
176// failure, or cancellation), the `Drop` implementation rolls back the reservation: it evicts the
177// key from the cache, queues a close request with Fxfs, and unblocks waiting tasks.
178struct PendingEntry {
179    cache: Arc<Cache>,
180    key: u64,
181}
182
183impl PendingEntry {
184    fn key(&self) -> u64 {
185        self.key
186    }
187
188    // Fulfills the pending entry by constructing the `CachedBlob`, atomically transitioning the
189    // cache entry from `Opening` to `Ready`, and notifying waiting tasks.
190    fn commit(self, vmo: zx::Vmo, size: u64) -> Arc<CachedBlob> {
191        let zero_children_receiver = ZeroChildrenReceiver { blob: OnceLock::new() };
192        let zero_children_registration =
193            fasync::EHandle::local().register_receiver(zero_children_receiver);
194
195        let (cached, event) = {
196            let mut state = self.cache.state.lock();
197            let (root_hash, event, verifier) = match state.blob_states_by_key.remove(&self.key) {
198                Some(BlobState::Opening { root_hash, event, verifier }) => {
199                    (root_hash, event, verifier)
200                }
201                _ => {
202                    unreachable!("Blob for key {} was unexpectedly not in Opening state", self.key)
203                }
204            };
205
206            let merkle_verifier = OnceLock::new();
207            if let Some(verifier) = verifier {
208                let _ = merkle_verifier.set(verifier);
209            }
210
211            let cached = Arc::new(CachedBlob {
212                vmo,
213                root_hash,
214                cache: Arc::downgrade(&self.cache),
215                vmo_key: self.key,
216                len: size,
217                registration: zero_children_registration,
218                merkle_verifier,
219            });
220
221            cached
222                .registration
223                .receiver()
224                .blob
225                .set(Arc::downgrade(&cached))
226                .expect("Failed to bind CachedBlob to ZeroChildrenReceiver: already initialized");
227
228            state.blob_states_by_key.insert(self.key, BlobState::Ready(cached.clone()));
229
230            (cached, event)
231        };
232
233        cached.wait_for_zero_children();
234        event.notify(usize::MAX);
235        cached
236    }
237}
238
239impl Drop for PendingEntry {
240    fn drop(&mut self) {
241        let event = {
242            let mut state = self.cache.state.lock();
243            match state.blob_states_by_key.get(&self.key) {
244                // If the blob committed successfully, there is nothing to roll back.
245                Some(BlobState::Ready(_)) => return,
246                Some(BlobState::Opening { .. }) => {
247                    if let Some(BlobState::Opening { root_hash, event, .. }) =
248                        state.blob_states_by_key.remove(&self.key)
249                    {
250                        state.keys_by_hash.remove(&root_hash);
251                        event
252                    } else {
253                        return;
254                    }
255                }
256                None => {
257                    log::error!(
258                        key = self.key;
259                        "PendingEntry::drop: key unexpectedly absent from cache"
260                    );
261                    return;
262                }
263            }
264        };
265
266        // Queue the close on Fxfs before waking waiters.
267        // Note: If `open()` was never called or returned an error, closing this key is safe
268        // because the block driver ignores `CloseBlob` for unknown keys. However, if `open()` was
269        // in-flight when cancelled, or if `open()` succeeded but subsequent operations failed, this
270        // close prevents leaking extent mappings.
271        let session = self.cache.mapping_session.clone();
272        let key = self.key;
273        self.cache.scope.spawn(async move {
274            let _ = session.close(key).await;
275        });
276
277        // Wake up all suspended tasks concurrently waiting on this blob.
278        event.notify(usize::MAX);
279    }
280}
281
282// The outcome of looking up or reserving an entry in `Cache`.
283enum CacheLookup {
284    // The VMO has been created and cached. Returns a child of it.
285    Ready(zx::Vmo),
286    // The VMO is currently being created by another request. Contains a listener to await
287    // completion.
288    Pending(event_listener::EventListener),
289    // The blob is absent. Provides a drop-safe pending entry to be used with creation. If the task
290    // crashes or returns an early error, the pending state will be cleared from the cache.
291    Missing(PendingEntry),
292}
293
294// Coordinates the concurrent creation, caching, and verification of pager-backed blobs.
295struct Cache {
296    state: Mutex<CacheState>,
297    mapping_session: fmapping::MappingSessionProxy,
298    pager: Arc<zx::Pager>,
299    delivery_vmo: zx::Vmo,
300    scope: fasync::ScopeHandle,
301}
302
303impl Cache {
304    fn new(
305        mapping_session: fmapping::MappingSessionProxy,
306        pager: Arc<zx::Pager>,
307        delivery_vmo: zx::Vmo,
308        scope: fasync::ScopeHandle,
309    ) -> Self {
310        Self {
311            state: Mutex::new(CacheState::default()),
312            mapping_session,
313            pager,
314            delivery_vmo,
315            scope,
316        }
317    }
318
319    fn delivery_vmo(&self) -> &zx::Vmo {
320        &self.delivery_vmo
321    }
322
323    fn get_by_key(&self, key: u64) -> Option<Arc<CachedBlob>> {
324        let state = self.state.lock();
325        match state.blob_states_by_key.get(&key) {
326            Some(BlobState::Ready(cached)) => Some(cached.clone()),
327            _ => None,
328        }
329    }
330
331    #[cfg(test)]
332    fn is_merkle_initialized(&self, key: u64) -> bool {
333        let state = self.state.lock();
334        match state.blob_states_by_key.get(&key) {
335            Some(BlobState::Opening { verifier, .. }) => verifier.is_some(),
336            Some(BlobState::Ready(blob)) => blob.merkle_verifier.get().is_some(),
337            None => false,
338        }
339    }
340
341    /// Looks up the blob in the cache and returns its current [`CacheLookup`] state:
342    ///
343    /// - [`CacheLookup::Ready`]: The blob is already cached. Returns a new read-only child VMO.
344    /// - [`CacheLookup::Pending`]: Another task is actively opening this blob. Returns an
345    ///   `EventListener` so the caller can await completion without duplicate work.
346    /// - [`CacheLookup::Missing`]: The blob is absent. Atomically reserves the key and slot
347    ///   by inserting `BlobState::Opening`, returning a [`PendingEntry`] for the caller to fulfill.
348    fn lookup(self: &Arc<Self>, identifier: &[u8; 32]) -> Result<CacheLookup, Error> {
349        let mut state = self.state.lock();
350        if let Some(&key) = state.keys_by_hash.get(identifier) {
351            match state.blob_states_by_key.get(&key) {
352                Some(BlobState::Ready(cached)) => {
353                    let child = cached
354                        .create_child()
355                        .map_err(|s| anyhow!("Failed to create child VMO: {s}"))?;
356                    return Ok(CacheLookup::Ready(child));
357                }
358                Some(BlobState::Opening { event, .. }) => {
359                    return Ok(CacheLookup::Pending(event.listen()));
360                }
361                None => unreachable!("keys_by_hash pointed to non-existent key {key}"),
362            }
363        }
364
365        let event = Arc::new(event_listener::Event::new());
366        let key = state.generate_key();
367        state.keys_by_hash.insert(*identifier, key);
368        state
369            .blob_states_by_key
370            .insert(key, BlobState::Opening { root_hash: *identifier, event, verifier: None });
371        Ok(CacheLookup::Missing(PendingEntry { cache: self.clone(), key }))
372    }
373
374    fn on_zero_children(&self, blob: &Arc<CachedBlob>) {
375        // On zero children, evict the blob from the cache and close its Fxfs mapping session.
376        // We must hold this lock before checking `num_children` so `lookup` cannot clone
377        // a child between this check and the following eviction, which would leave that client
378        // holding a VMO whose storage mappings have already been closed.
379        let mut state = self.state.lock();
380        if let Ok(info) = blob.vmo.info() {
381            if info.num_children > 0 {
382                // A concurrent client acquired a child VMO before we processed this packet.
383                // Resume watching for the next ZERO_CHILDREN signal.
384                blob.wait_for_zero_children();
385                return;
386            }
387        }
388        // Evict from `keys_by_hash`.
389        if let hash_map::Entry::Occupied(entry) = state.keys_by_hash.entry(blob.root_hash) {
390            if *entry.get() == blob.vmo_key {
391                entry.remove();
392            }
393        }
394        // Evict from `blob_states_by_key`.
395        if let hash_map::Entry::Occupied(entry) = state.blob_states_by_key.entry(blob.vmo_key) {
396            if let BlobState::Ready(cached) = entry.get() {
397                if Arc::ptr_eq(cached, blob) {
398                    entry.remove();
399                }
400            }
401        }
402    }
403
404    fn create_merkle_verifier(
405        root_hash: &[u8; 32],
406        vmo_key: u64,
407        hashes: Box<[Hash]>,
408    ) -> Result<ReadSizedMerkleVerifier, Error> {
409        // For blobs <= 8192 bytes (a single block), the Merkle tree consists of only a single leaf,
410        // which is identical to the root hash itself. Fxfs omits storing leaf hashes on disk for
411        // single-block blobs to save space, resulting in empty leaf metadata read by the block
412        // server. In that case, synthesise the single leaf from the blob's Merkle root hash.
413        let hashes = if hashes.is_empty() { Box::new([Hash::from(*root_hash)]) } else { hashes };
414
415        match MerkleVerifier::new(Hash::from(*root_hash), hashes) {
416            Ok(verifier) => ReadSizedMerkleVerifier::new(verifier, delivery::DELIVERY_DATA_SIZE)
417                .map_err(|e| anyhow!("Failed to create ReadSizedMerkleVerifier: {e:?}")),
418            Err(e) => {
419                bail!("Failed to verify merkle leaves for key {vmo_key}: {e:?}");
420            }
421        }
422    }
423
424    fn report_pager_failure(&self, vmo: &zx::Vmo, range: Range<u64>, status: zx::Status) {
425        // Map to the set of statuses accepted by ZX_PAGER_OP_FAIL which are ZX_ERR_IO,
426        // ZX_ERR_IO_DATA_INTEGRITY, ZX_ERR_BAD_STATE, ZX_ERR_NO_SPACE, and ZX_ERR_BUFFER_TOO_SMALL.
427        let pager_status = match status {
428            zx::Status::IO_DATA_INTEGRITY => zx::Status::IO_DATA_INTEGRITY,
429            zx::Status::NO_SPACE => zx::Status::NO_SPACE,
430            zx::Status::BUFFER_TOO_SMALL | zx::Status::FILE_BIG => zx::Status::BUFFER_TOO_SMALL,
431            zx::Status::IO
432            | zx::Status::IO_DATA_LOSS
433            | zx::Status::IO_INVALID
434            | zx::Status::IO_MISSED_DEADLINE
435            | zx::Status::IO_NOT_PRESENT
436            | zx::Status::IO_OVERRUN
437            | zx::Status::IO_REFUSED
438            | zx::Status::PEER_CLOSED => zx::Status::IO,
439            _ => zx::Status::BAD_STATE,
440        };
441
442        // `ZX_PAGER_OP_FAIL` resolves pending page requests overlapping the specified range with
443        // the failure status. Pages that were already supplied remain accessible.
444        if let Err(error) = self.pager.op_range(zx::PagerOp::Fail(pager_status), vmo, range) {
445            log::error!(error:?; "Failed to report pager failure to kernel");
446        }
447    }
448}
449
450impl delivery::DeliveryQueueProvider for Cache {
451    fn deliver_pages(
452        &self,
453        key: u64,
454        target_offset: u64,
455        unverified_pages: UnverifiedPages<'_>,
456    ) -> Result<(), Error> {
457        let cached_blob = self
458            .get_by_key(key)
459            .ok_or_else(|| anyhow!("Unknown or expired key {key} in deliver_pages"))?;
460
461        let page_size = zx::system_get_page_size() as u64;
462        let page_aligned_size = cached_blob.len.div_ceil(page_size) * page_size;
463
464        let verifier = match cached_blob.merkle_verifier.get() {
465            Some(v) => v,
466            None => {
467                // If data arrives before `RegisterBlob` (or if registration failed), the blob has
468                // no cryptographic metadata, so no pages can ever be verified or supplied. Fail the
469                // entire page-aligned range of the VMO with `BAD_STATE` to unblock any waiting
470                // client threads on this blob.
471                self.report_pager_failure(
472                    &cached_blob.vmo,
473                    0..page_aligned_size,
474                    zx::Status::BAD_STATE,
475                );
476                bail!("Data received for uninitialized blob {}", key);
477            }
478        };
479
480        let chunk_len = unverified_pages.len_in_bytes() as u64;
481
482        // `zx_pager_supply_pages` and `zx_pager_op_range(Fail)` require the range to be within the
483        // page-aligned capacity of the VMO. If the range extends past the end of the VMO, the
484        // kernel fails the call with `ZX_ERR_OUT_OF_RANGE`. Clamp the range to the VMO's
485        // page-aligned limit.
486        let bounded_len = std::cmp::min(chunk_len, page_aligned_size.saturating_sub(target_offset));
487        let chunk_range = target_offset..target_offset + bounded_len;
488
489        // For the final chunk of an unaligned blob, the buffer passed to `verify_aligned` is
490        // page-aligned and zero-padded, but `verify_aligned` expects the remaining unaligned length
491        // of the valid data.
492        let remaining_in_blob = cached_blob.len.saturating_sub(target_offset);
493        let unaligned_len = std::cmp::min(chunk_len, remaining_in_blob) as usize;
494        if let Err(e) = verifier.verify_aligned(
495            target_offset as usize,
496            unverified_pages.as_ptr_byte_slice(),
497            unaligned_len,
498        ) {
499            self.report_pager_failure(&cached_blob.vmo, chunk_range, zx::Status::IO_DATA_INTEGRITY);
500            bail!("Failed to verify payload for blob {}: {:?}", key, e);
501        }
502
503        if let Err(e) = self.pager.supply_pages(
504            &cached_blob.vmo,
505            chunk_range.clone(),
506            &self.delivery_vmo,
507            unverified_pages.delivery_offset(),
508        ) {
509            self.report_pager_failure(&cached_blob.vmo, chunk_range, e);
510            bail!("Failed to supply pages: {:?}", e);
511        }
512
513        Ok(())
514    }
515
516    fn register_blob(&self, key: u64, leaf_data: PtrByteSlice<'_>) -> Result<(), Error> {
517        if leaf_data.len() % HASH_SIZE != 0 {
518            bail!("RegisterBlob invalid leaf length must be a multiple of HASH_SIZE");
519        }
520
521        let hashes: Vec<Hash> = (0..(leaf_data.len() / HASH_SIZE))
522            .map(|i| {
523                let chunk = leaf_data.subslice(i * HASH_SIZE..(i + 1) * HASH_SIZE);
524                let mut hash_bytes = [0u8; HASH_SIZE];
525                chunk.copy_to_slice(&mut hash_bytes);
526                Hash::from(hash_bytes)
527            })
528            .collect();
529
530        let mut state = self.state.lock();
531        match state.blob_states_by_key.get_mut(&key) {
532            Some(BlobState::Opening { root_hash, verifier, .. }) => {
533                if verifier.is_some() {
534                    log::warn!(key;
535                        "RegisterBlob received for blob but it is already initialized"
536                    );
537                    return Ok(());
538                }
539                *verifier =
540                    Some(Self::create_merkle_verifier(root_hash, key, hashes.into_boxed_slice())?);
541                Ok(())
542            }
543            Some(BlobState::Ready(cached)) => {
544                if cached.merkle_verifier.get().is_some() {
545                    log::warn!(key;
546                        "RegisterBlob received for blob but it is already initialized"
547                    );
548                    return Ok(());
549                }
550                let verifier = Self::create_merkle_verifier(
551                    &cached.root_hash,
552                    key,
553                    hashes.into_boxed_slice(),
554                )?;
555                let _ = cached.merkle_verifier.set(verifier);
556                Ok(())
557            }
558            None => {
559                bail!("RegisterBlob received for unknown or expired key {key}");
560            }
561        }
562    }
563}
564
565/// Coordinates between FxBlob, the block driver, and the kernel to manage pager-backed blobs such
566/// that page requests can be handled at the driver instead of the filesystem layer.
567///
568/// `BlobPagerAndVerifier` uses a dedicated Pager to create pager-backed VMOs for blobs. Kernel page
569/// faults on these VMOs are routed directly to the block driver, which reads and decompresses the
570/// raw data. The driver then passes this data back to `BlobPagerAndVerifier` for cryptographic
571/// verification and page fault resolution.
572pub struct BlobPagerAndVerifier {
573    // Ties the lifetime of the background delivery thread to the `BlobPagerAndVerifier`.
574    _delivery_processor: delivery::DeliveryQueueProcessor,
575    // Ties the lifetime of the block driver's mapper session to the `BlobPagerAndVerifier`.
576    _mapper_session: fblock::MapperSessionProxy,
577    port: zx::Port,
578    pager: Arc<zx::Pager>,
579    cache: Arc<Cache>,
580    _scope: fasync::Scope,
581}
582
583impl BlobPagerAndVerifier {
584    /// Returns a reference to the shared delivery queue VMO.
585    pub fn delivery_vmo(&self) -> &zx::Vmo {
586        self.cache.delivery_vmo()
587    }
588
589    /// Creates a new BlobPagerAndVerifier.
590    ///
591    /// Establishes a mapping_provider session with FxBlob to receive a shared VMO that will be used
592    /// to communicate opened blob extent mappings as well as any closed blobs. Then, opens a
593    /// mapper session with the block driver, forwarding the mapping VMO alongside a Zircon port
594    /// (for the driver to receive page faults) and a delivery queue VMO (for the driver to emit
595    /// unverified data for verification).
596    pub async fn new(
597        mapping_provider: &fmapping::MappingProviderProxy,
598        mapper: &fblock::MapperProxy,
599    ) -> Result<Self, Error> {
600        // Establish a mapping_provider session with Fxfs.
601        let (mapping_session, session_server) =
602            fidl::endpoints::create_proxy::<fmapping::MappingSessionMarker>();
603        let mapping_vmo = mapping_provider
604            .open_session(session_server)
605            .await
606            .context("FIDL error calling MappingProvider.OpenSession")?
607            .map_err(|e| anyhow!("MappingProvider open_session failed: {e:?}"))?;
608
609        // Establish a mapper session with the block driver.
610        let port = zx::Port::create();
611        let delivery_queue = zx::Vmo::create(mapping::DELIVERY_VMO_SIZE)
612            .context("Failed to create delivery queue VMO")?;
613        let delivery_vmo: zx::Vmo = delivery_queue
614            .duplicate_handle(zx::Rights::SAME_RIGHTS)
615            .context("Failed to duplicate delivery queue VMO")?;
616        let delivery_queue_dup = delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS)?;
617        let receiver = vmo_fifo::Receiver::<mapping::RawDeliveryCommand>::new(
618            delivery_queue,
619            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
620        )
621        .context("Failed to create delivery queue receiver")?;
622        let pager = Arc::new(
623            zx::Pager::create(zx::PagerOptions::empty()).context("Failed to create pager")?,
624        );
625        let port_dup = port.duplicate_handle(zx::Rights::SAME_RIGHTS)?;
626        let (mapper_session, mapper_session_server) =
627            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
628        mapper
629            .open_session(
630                mapper_session_server,
631                mapping_vmo,
632                Some(port_dup),
633                Some(delivery_queue_dup),
634            )
635            .await
636            .context("FIDL error calling Mapper.OpenSession")?
637            .map_err(|e| anyhow!("Mapper.OpenSession failed: {e:?}"))?;
638
639        let scope = fasync::Scope::new();
640        let cache = Arc::new(Cache::new(
641            mapping_session,
642            pager.clone(),
643            delivery_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)?,
644            scope.to_handle(),
645        ));
646        let _delivery_processor =
647            delivery::DeliveryQueueProcessor::spawn(receiver, cache.clone(), delivery_vmo)?;
648
649        Ok(Self {
650            _delivery_processor,
651            _mapper_session: mapper_session,
652            port,
653            pager,
654            cache,
655            _scope: scope,
656        })
657    }
658
659    /// Create pager owned VMO for the blob identified by its Merkle Root Hash.
660    pub async fn create_vmo(&self, identifier: &[u8; 32]) -> Result<zx::Vmo, Error> {
661        loop {
662            let entry = match self.cache.lookup(identifier)? {
663                CacheLookup::Ready(child) => return Ok(child),
664                CacheLookup::Pending(listener) => {
665                    listener.await;
666                    continue; // Re-evaluate cache now that the blocking event triggered
667                }
668                CacheLookup::Missing(entry) => entry,
669            };
670
671            let key = entry.key();
672
673            // Ask Fxfs to register the blob and write its extent mappings into the shared mapping
674            // VMO, using the client-allocated key.
675            let size = self
676                .cache
677                .mapping_session
678                .open(key, identifier)
679                .await
680                .context("FIDL error calling MappingSession.Open")?
681                .map_err(|e| anyhow!("MappingSession.Open failed for blob: {e:?}"))?;
682
683            let vmo = self.pager.create_vmo(zx::VmoOptions::empty(), &self.port, key, size)?;
684
685            // Create the initial child. We vend children of this VMO to clients so we can track
686            // when all children are dropped to evict the blob from cache.
687            let first_child = vmo
688                .create_child(zx::VmoChildOptions::REFERENCE | zx::VmoChildOptions::NO_WRITE, 0, 0)
689                .map_err(|s| anyhow!("Failed to create child VMO: {}", s))?;
690
691            entry.commit(vmo, size);
692
693            return Ok(first_child);
694        }
695    }
696}
697
698#[cfg(test)]
699mod tests {
700    use super::*;
701    use crate::delivery::DELIVERY_DATA_SIZE;
702    use futures::TryStreamExt;
703    use futures::channel::oneshot;
704    use mapping::{DeliveryCommand, RawDeliveryCommand};
705    use vmo_fifo::SyncSender;
706
707    const TEST_BLOB_SIZE: u64 = (DELIVERY_DATA_SIZE * 2) as u64;
708    const TEST_VMO_KEY: u64 = 1;
709
710    // Used for testing BlobPagerAndVerifier interactions.
711    // We arbitrarily allocate a 32KiB buffer so the payload spans multiple VMO pages.
712    //
713    // This test environment assumes a single blob file, identified by the `valid_root` hash.
714    struct TestEnv {
715        pager_and_verifier: Arc<BlobPagerAndVerifier>,
716        mapping_task: fasync::Task<()>,
717        mapper_task: fasync::Task<()>,
718        close_signal: Option<oneshot::Receiver<()>>,
719        delivery_vmo: zx::Vmo,
720        pub valid_root: [u8; 32],
721        pub valid_leaves: Vec<u8>,
722        pub blob_data: Vec<u8>,
723    }
724
725    impl TestEnv {
726        async fn new(blob_size: u64) -> Self {
727            let blob_data = vec![0x42u8; blob_size as usize];
728            let (root, leaf_hashes) =
729                fuchsia_merkle::MerkleRootBuilder::new(Vec::new()).complete(&blob_data);
730            let expected_hash: [u8; 32] = root.into();
731
732            let mut flat_leaves = Vec::new();
733            for hash in &leaf_hashes {
734                flat_leaves.extend_from_slice(hash.as_bytes());
735            }
736
737            let (mapping_proxy, mut mapping_stream) =
738                fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
739
740            let (close_tx, close_rx) = oneshot::channel();
741            let mapping_task = fasync::Task::spawn(async move {
742                let mut close_tx = Some(close_tx);
743                if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
744                    mapping_stream.try_next().await.expect("try_next failed")
745                {
746                    let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
747                        .expect("zx::Vmo::create failed");
748                    responder.send(Ok(mapping_vmo)).expect("send failed");
749
750                    let mut session_stream = session.into_stream();
751                    let mut open_calls = 0;
752                    while let Some(request) =
753                        session_stream.try_next().await.expect("try_next failed")
754                    {
755                        match request {
756                            fmapping::MappingSessionRequest::Open {
757                                key,
758                                identifier,
759                                responder,
760                            } => {
761                                assert_eq!(identifier, expected_hash);
762                                assert_eq!(key, TEST_VMO_KEY);
763                                open_calls += 1;
764                                assert_eq!(
765                                    open_calls, 1,
766                                    "Open should only be called once per identifier"
767                                );
768                                responder.send(Ok(blob_size)).expect("send failed");
769                            }
770                            fmapping::MappingSessionRequest::Close { key, responder } => {
771                                assert_eq!(key, TEST_VMO_KEY);
772                                responder.send(Ok(())).expect("send failed");
773                                if let Some(tx) = close_tx.take() {
774                                    let _ = tx.send(());
775                                }
776                            }
777                            _ => {}
778                        }
779                    }
780                }
781            });
782
783            let (mapper_proxy, mut mapper_stream) =
784                fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
785
786            let (tx_delivery, rx_delivery) = oneshot::channel();
787            let mapper_task = fasync::Task::spawn(async move {
788                if let Some(fblock::MapperRequest::OpenSession {
789                    delivery_queue, responder, ..
790                }) = mapper_stream.try_next().await.expect("try_next failed")
791                {
792                    tx_delivery.send(delivery_queue).expect("send failed");
793                    responder.send(Ok(())).expect("send failed");
794                }
795            });
796
797            Self {
798                pager_and_verifier: Arc::new(
799                    BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
800                        .await
801                        .expect("BlobPagerAndVerifier::new failed"),
802                ),
803                mapping_task,
804                mapper_task,
805                close_signal: Some(close_rx),
806                delivery_vmo: rx_delivery
807                    .await
808                    .expect("rx_delivery wait failed")
809                    .expect("delivery_vmo was None"),
810                valid_root: expected_hash,
811                valid_leaves: flat_leaves,
812                blob_data,
813            }
814        }
815
816        async fn teardown(self) {
817            drop(self.pager_and_verifier);
818            self.mapping_task.await;
819            self.mapper_task.await;
820        }
821    }
822
823    struct ExpectedReadThread {
824        started_rx: Option<oneshot::Receiver<()>>,
825        done_rx: oneshot::Receiver<()>,
826        handle: std::thread::JoinHandle<()>,
827    }
828
829    impl ExpectedReadThread {
830        async fn wait_started(&mut self) {
831            if let Some(rx) = self.started_rx.take() {
832                rx.await.expect("reader thread failed to start");
833            }
834        }
835
836        async fn wait_and_verify(self) {
837            self.done_rx.await.expect("reader thread panicked or hung");
838            self.handle.join().expect("join failed");
839        }
840    }
841
842    fn spawn_reader_expect_error(
843        vmo: &zx::Vmo,
844        offset: u64,
845        size: usize,
846        expected_status: zx::Status,
847    ) -> ExpectedReadThread {
848        let vmo_clone =
849            vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate_handle failed");
850        let (started_tx, started_rx) = oneshot::channel();
851        let (done_tx, done_rx) = oneshot::channel();
852
853        let handle = std::thread::spawn(move || {
854            let mut buf = vec![0u8; size];
855            let _ = started_tx.send(());
856            let err = vmo_clone.read(&mut buf, offset).expect_err("read should fail");
857            assert_eq!(err, expected_status);
858            let _ = done_tx.send(());
859        });
860
861        ExpectedReadThread { started_rx: Some(started_rx), done_rx, handle }
862    }
863
864    #[fuchsia::test]
865    async fn test_create_vmo() {
866        let env = TestEnv::new(TEST_BLOB_SIZE).await;
867        let vmo =
868            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("Failed to create VMO");
869        let size = vmo.get_size().expect("get_size failed");
870
871        // Pager VMO sizes are rounded up to the nearest page boundary.
872        let page_size = zx::system_get_page_size() as u64;
873        let expected_pages = (TEST_BLOB_SIZE + page_size - 1) / page_size;
874        assert_eq!(size, expected_pages * page_size);
875
876        // Explicitly drop the VMO so `ZX_VMO_ZERO_CHILDREN` fires and the mock mapping
877        // server receives the `.close()` IPC. Otherwise `teardown` will deadlock waiting
878        // for the server loop to exit!
879        drop(vmo);
880
881        env.teardown().await;
882    }
883
884    #[fuchsia::test]
885    async fn test_create_vmo_concurrent_access() {
886        let env = TestEnv::new(TEST_BLOB_SIZE).await;
887        let mut futures = vec![];
888        for _ in 0..10 {
889            let verifier = env.pager_and_verifier.clone();
890            futures.push(fasync::Task::spawn(async move {
891                verifier.create_vmo(&env.valid_root).await.expect("Failed to create VMO")
892            }));
893        }
894        let vmos = futures::future::join_all(futures).await;
895
896        drop(vmos);
897        env.teardown().await;
898    }
899
900    #[fuchsia::test]
901    async fn test_zero_children_eviction() {
902        let mut env = TestEnv::new(TEST_BLOB_SIZE).await;
903
904        let child_vmo =
905            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
906
907        {
908            let state = env.pager_and_verifier.cache.state.lock();
909            assert_eq!(state.keys_by_hash.get(&env.valid_root), Some(&TEST_VMO_KEY));
910            assert!(matches!(
911                state.blob_states_by_key.get(&TEST_VMO_KEY),
912                Some(BlobState::Ready(_))
913            ));
914        }
915
916        drop(child_vmo);
917
918        // wait for the packet receiver to observe the ZERO_CHILDREN signal and dispatch
919        // `mapping_session.close()` to the mapping_server.
920        env.close_signal.take().expect("Missing signal").await.expect("Failed to close");
921
922        {
923            let state = env.pager_and_verifier.cache.state.lock();
924            assert!(state.keys_by_hash.get(&env.valid_root).is_none());
925            assert!(state.blob_states_by_key.get(&TEST_VMO_KEY).is_none());
926        }
927
928        env.teardown().await;
929    }
930
931    #[fuchsia::test]
932    async fn test_create_vmo_multiple_distinct_blobs() {
933        use futures::StreamExt;
934
935        let (mapping_proxy, mut mapping_stream) =
936            fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
937        let (mapper_proxy, mut mapper_stream) =
938            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
939
940        let (close_tx, mut close_rx) = futures::channel::mpsc::unbounded();
941        let mapping_task = fasync::Task::spawn(async move {
942            if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
943                mapping_stream.try_next().await.expect("try_next failed")
944            {
945                let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
946                    .expect("zx::Vmo::create failed");
947                responder.send(Ok(mapping_vmo)).expect("send failed");
948
949                let mut session_stream = session.into_stream();
950                while let Some(request) = session_stream.try_next().await.expect("try_next failed")
951                {
952                    match request {
953                        fmapping::MappingSessionRequest::Open { responder, .. } => {
954                            responder.send(Ok(8192)).expect("send failed");
955                        }
956                        fmapping::MappingSessionRequest::Close { key, responder } => {
957                            responder.send(Ok(())).expect("send failed");
958                            let _ = close_tx.unbounded_send(key);
959                        }
960                        _ => {}
961                    }
962                }
963            }
964        });
965
966        let mapper_task = fasync::Task::spawn(async move {
967            if let Some(fblock::MapperRequest::OpenSession { responder, .. }) =
968                mapper_stream.try_next().await.expect("try_next failed")
969            {
970                responder.send(Ok(())).expect("send failed");
971            }
972        });
973
974        let pager_and_verifier = Arc::new(
975            BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
976                .await
977                .expect("BlobPagerAndVerifier::new failed"),
978        );
979
980        let hash1: [u8; 32] = [0x11; 32];
981        let hash2: [u8; 32] = [0x22; 32];
982        let hash3: [u8; 32] = [0x33; 32];
983
984        let vmo1 = pager_and_verifier.create_vmo(&hash1).await.expect("create_vmo 1 failed");
985        let vmo2 = pager_and_verifier.create_vmo(&hash2).await.expect("create_vmo 2 failed");
986        let vmo3 = pager_and_verifier.create_vmo(&hash3).await.expect("create_vmo 3 failed");
987
988        let (key1, key2, key3) = {
989            let state = pager_and_verifier.cache.state.lock();
990            let k1 = *state.keys_by_hash.get(&hash1).expect("hash1 missing");
991            let k2 = *state.keys_by_hash.get(&hash2).expect("hash2 missing");
992            let k3 = *state.keys_by_hash.get(&hash3).expect("hash3 missing");
993            assert_ne!(k1, k2);
994            assert_ne!(k2, k3);
995            assert_ne!(k1, k3);
996            assert!(matches!(state.blob_states_by_key.get(&k1), Some(BlobState::Ready(_))));
997            assert!(matches!(state.blob_states_by_key.get(&k2), Some(BlobState::Ready(_))));
998            assert!(matches!(state.blob_states_by_key.get(&k3), Some(BlobState::Ready(_))));
999            (k1, k2, k3)
1000        };
1001
1002        // Evict blob 2 by dropping its only child VMO
1003        drop(vmo2);
1004        let closed_key = close_rx.next().await.expect("expected close signal");
1005        assert_eq!(closed_key, key2);
1006
1007        // Verify blob 2 was evicted while blob 1 and 3 remain cached
1008        {
1009            let state = pager_and_verifier.cache.state.lock();
1010            assert!(state.keys_by_hash.get(&hash2).is_none());
1011            assert!(state.blob_states_by_key.get(&key2).is_none());
1012            assert_eq!(state.keys_by_hash.get(&hash1), Some(&key1));
1013            assert_eq!(state.keys_by_hash.get(&hash3), Some(&key3));
1014            assert!(matches!(state.blob_states_by_key.get(&key1), Some(BlobState::Ready(_))));
1015            assert!(matches!(state.blob_states_by_key.get(&key3), Some(BlobState::Ready(_))));
1016        }
1017
1018        // Re-open blob 2: it should generate a new key and open successfully
1019        let vmo2_new = pager_and_verifier.create_vmo(&hash2).await.expect("re-create vmo 2 failed");
1020        {
1021            let state = pager_and_verifier.cache.state.lock();
1022            let k = *state.keys_by_hash.get(&hash2).expect("hash2 missing after re-open");
1023            assert_ne!(k, key2);
1024            assert!(matches!(state.blob_states_by_key.get(&k), Some(BlobState::Ready(_))));
1025        }
1026
1027        drop(vmo1);
1028        drop(vmo3);
1029        drop(vmo2_new);
1030
1031        drop(pager_and_verifier);
1032        mapping_task.await;
1033        mapper_task.await;
1034    }
1035
1036    #[fuchsia::test]
1037    async fn test_cache_is_cleared_when_create_vmo_fails() {
1038        let (mapping_proxy, mut mapping_stream) =
1039            fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
1040        let (mapper_proxy, mut mapper_stream) =
1041            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
1042
1043        let (mock_opened_tx, mock_opened_rx) = oneshot::channel();
1044        let mapping_task = fasync::Task::spawn(async move {
1045            let mut mock_opened_tx = Some(mock_opened_tx);
1046            if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
1047                mapping_stream.try_next().await.expect("try_next failed")
1048            {
1049                let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
1050                    .expect("zx::Vmo::create failed");
1051                responder.send(Ok(mapping_vmo)).expect("send failed");
1052                let mut session_stream = session.into_stream();
1053
1054                if let Some(fmapping::MappingSessionRequest::Open {
1055                    responder: _responder, ..
1056                }) = session_stream.try_next().await.expect("try_next failed")
1057                {
1058                    // Signal the primary test task that the session is inside Open
1059                    if let Some(tx) = mock_opened_tx.take() {
1060                        let _ = tx.send(());
1061                    }
1062
1063                    // Wait for test task to simulate failing create_vmo
1064                    let () = std::future::pending().await;
1065                }
1066            }
1067        });
1068
1069        let mapper_task = fasync::Task::spawn(async move {
1070            if let Some(fblock::MapperRequest::OpenSession { responder, .. }) =
1071                mapper_stream.try_next().await.expect("try_next failed")
1072            {
1073                responder.send(Ok(())).expect("send failed");
1074            }
1075        });
1076
1077        let pager_and_verifer = Arc::new(
1078            BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
1079                .await
1080                .expect("BlobPagerAndVerifier::new failed"),
1081        );
1082
1083        let blob_data = vec![0x42u8; TEST_BLOB_SIZE as usize];
1084        let (root, _) = fuchsia_merkle::MerkleRootBuilder::new(Vec::new()).complete(&blob_data);
1085        let hash_val: [u8; 32] = root.into();
1086
1087        let verifier_clone = pager_and_verifer.clone();
1088
1089        let identifier = hash_val;
1090        let (abortable_future, abort_handle) = futures::future::abortable(async move {
1091            let _ = verifier_clone.create_vmo(&identifier).await;
1092        });
1093
1094        let identifier = &hash_val;
1095        let primary_task = fasync::Task::spawn(abortable_future);
1096
1097        // Wait until the mock mapping session catches the `Open` request and check that the cache
1098        // state is updated.
1099        mock_opened_rx.await.expect("Open exited abruptly");
1100        {
1101            let state = pager_and_verifer.cache.state.lock();
1102            assert_eq!(state.keys_by_hash.get(identifier), Some(&1));
1103            assert!(matches!(state.blob_states_by_key.get(&1), Some(BlobState::Opening { .. })));
1104        }
1105
1106        // Simulate a failed create_vmo (future abort drops execution context)
1107        abort_handle.abort();
1108        let _ = primary_task.await;
1109
1110        // The cache should be cleared after due to PendingEntry being dropped
1111        {
1112            let state = pager_and_verifer.cache.state.lock();
1113            assert!(state.keys_by_hash.get(identifier).is_none());
1114            assert!(state.blob_states_by_key.get(&1).is_none());
1115        }
1116
1117        drop(pager_and_verifer);
1118        drop(mapping_task);
1119        drop(mapper_task);
1120    }
1121
1122    #[fuchsia::test]
1123    async fn test_create_vmo_error_clears_cache_and_unblocks_waiters() {
1124        let (mapping_proxy, mut mapping_stream) =
1125            fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
1126        let (mapper_proxy, mut mapper_stream) =
1127            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
1128
1129        let (open_received_tx, open_received_rx) = oneshot::channel();
1130        let (reply_tx, reply_rx) = oneshot::channel();
1131
1132        let mapping_task = fasync::Task::spawn(async move {
1133            if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
1134                mapping_stream.try_next().await.expect("try_next failed")
1135            {
1136                let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
1137                    .expect("zx::Vmo::create failed");
1138                responder.send(Ok(mapping_vmo)).expect("send failed");
1139                let mut session_stream = session.into_stream();
1140
1141                if let Some(fmapping::MappingSessionRequest::Open { responder, .. }) =
1142                    session_stream.try_next().await.expect("try_next failed")
1143                {
1144                    // Signal the test that the primary caller has entered Open and is Pending.
1145                    open_received_tx.send(()).expect("send open_received failed");
1146
1147                    // Await permission to reply with the error.
1148                    reply_rx.await.expect("reply_rx failed");
1149                    responder.send(Err(zx::Status::NOT_FOUND.into_raw())).expect("send failed");
1150                }
1151                // Channel closes when mapping_task finishes, causing subsequent requests to fail.
1152            }
1153        });
1154
1155        let mapper_task = fasync::Task::spawn(async move {
1156            if let Some(fblock::MapperRequest::OpenSession { responder, .. }) =
1157                mapper_stream.try_next().await.expect("try_next failed")
1158            {
1159                responder.send(Ok(())).expect("send failed");
1160            }
1161        });
1162
1163        let pager_and_verifier = Arc::new(
1164            BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
1165                .await
1166                .expect("BlobPagerAndVerifier::new failed"),
1167        );
1168
1169        let hash: [u8; 32] = [0x55; 32];
1170
1171        // Primary caller starts create_vmo.
1172        let verifier1 = pager_and_verifier.clone();
1173        let primary_future = fasync::Task::spawn(async move { verifier1.create_vmo(&hash).await });
1174
1175        // Wait until the primary caller is handling OpenSession (cache is Pending).
1176        open_received_rx.await.expect("open_received_rx failed");
1177        {
1178            let state = pager_and_verifier.cache.state.lock();
1179            assert!(matches!(
1180                state.keys_by_hash.get(&hash).and_then(|k| state.blob_states_by_key.get(k)),
1181                Some(BlobState::Opening { .. })
1182            ));
1183        }
1184
1185        // A second caller arrives while the primary creation is still pending.
1186        let verifier2 = pager_and_verifier.clone();
1187        let secondary_future =
1188            fasync::Task::spawn(async move { verifier2.create_vmo(&hash).await });
1189
1190        // Allow Fxfs to reply with an error to the primary caller.
1191        reply_tx.send(()).expect("reply_tx send failed");
1192
1193        // Both callers must return an error and complete without hanging on abandoned state.
1194        assert!(primary_future.await.is_err());
1195        assert!(secondary_future.await.is_err());
1196
1197        // Cache should be empty.
1198        {
1199            let state = pager_and_verifier.cache.state.lock();
1200            assert!(state.keys_by_hash.get(&hash).is_none());
1201            assert!(state.blob_states_by_key.is_empty());
1202        }
1203
1204        drop(pager_and_verifier);
1205        drop(mapping_task);
1206        drop(mapper_task);
1207    }
1208
1209    #[fuchsia::test]
1210    async fn test_register_blob_arrives_before_open_completes() {
1211        // Test that if `RegisterBlob` arrives while `session.open().await` is still pending
1212        // (the arrival race), the blob's Merkle tree is retained and page reads succeed once
1213        // `create_vmo` completes.
1214        let (mapping_proxy, mut mapping_stream) =
1215            fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
1216        let (mapper_proxy, mut mapper_stream) =
1217            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
1218
1219        let blob_data = vec![0x42u8; TEST_BLOB_SIZE as usize];
1220        let (root, leaf_hashes) =
1221            fuchsia_merkle::MerkleRootBuilder::new(Vec::new()).complete(&blob_data);
1222        let expected_hash: [u8; 32] = root.into();
1223
1224        let mut flat_leaves = Vec::new();
1225        for hash in &leaf_hashes {
1226            flat_leaves.extend_from_slice(hash.as_bytes());
1227        }
1228
1229        let (tx_delivery, rx_delivery) = oneshot::channel();
1230        let (open_key_tx, open_key_rx) = oneshot::channel();
1231        let (resume_open_tx, resume_open_rx) = oneshot::channel();
1232        let (close_tx, close_rx) = oneshot::channel();
1233
1234        let mapping_task = fasync::Task::spawn(async move {
1235            let mut close_tx = Some(close_tx);
1236            let mut open_key_tx = Some(open_key_tx);
1237            let mut resume_open_rx = Some(resume_open_rx);
1238            if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
1239                mapping_stream.try_next().await.expect("try_next failed")
1240            {
1241                let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
1242                    .expect("zx::Vmo::create failed");
1243                responder.send(Ok(mapping_vmo)).expect("send failed");
1244
1245                let mut session_stream = session.into_stream();
1246                while let Some(request) = session_stream.try_next().await.expect("try_next failed")
1247                {
1248                    match request {
1249                        fmapping::MappingSessionRequest::Open { key, identifier, responder } => {
1250                            assert_eq!(identifier, expected_hash);
1251                            // Send key to test thread so it can deliver leaves before Open returns
1252                            if let Some(tx) = open_key_tx.take() {
1253                                tx.send(key).expect("open_key_tx failed");
1254                            }
1255                            if let Some(rx) = resume_open_rx.take() {
1256                                rx.await.expect("resume_open_rx failed");
1257                            }
1258                            responder.send(Ok(TEST_BLOB_SIZE)).expect("send failed");
1259                        }
1260                        fmapping::MappingSessionRequest::Close { key: _, responder } => {
1261                            responder.send(Ok(())).expect("send failed");
1262                            if let Some(tx) = close_tx.take() {
1263                                let _ = tx.send(());
1264                            }
1265                        }
1266                        _ => {}
1267                    }
1268                }
1269            }
1270        });
1271
1272        let mapper_task = fasync::Task::spawn(async move {
1273            if let Some(fblock::MapperRequest::OpenSession { delivery_queue, responder, .. }) =
1274                mapper_stream.try_next().await.expect("try_next failed")
1275            {
1276                tx_delivery.send(delivery_queue).expect("send failed");
1277                responder.send(Ok(())).expect("send failed");
1278            }
1279        });
1280
1281        let pager_and_verifier = Arc::new(
1282            BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
1283                .await
1284                .expect("BlobPagerAndVerifier::new failed"),
1285        );
1286
1287        let delivery_vmo =
1288            rx_delivery.await.expect("rx_delivery wait failed").expect("delivery_vmo was None");
1289        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1290            delivery_vmo
1291                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1292                .expect("duplicate_handle failed"),
1293            zx::system_get_page_size() as usize,
1294            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1295        )
1296        .expect("SyncSender::new failed");
1297
1298        let verifier_clone = pager_and_verifier.clone();
1299        let create_task =
1300            fasync::Task::spawn(async move { verifier_clone.create_vmo(&expected_hash).await });
1301
1302        // Wait until `open()` is called and get the client-allocated key
1303        let key = open_key_rx.await.expect("open_key_rx failed");
1304
1305        // While `open()` is still pending, send `RegisterBlob`
1306        let mut payload =
1307            sender.reserve_payload(flat_leaves.len()).expect("reserve_payload failed");
1308        payload.data().copy_from_slice(&flat_leaves);
1309        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1310            key: key as u64,
1311            offset: payload.offset(),
1312            length: flat_leaves.len() as u32,
1313        }
1314        .into();
1315        payload.commit(raw_cmd).expect("commit failed");
1316
1317        // Wait until delivery thread processes `RegisterBlob` on the pending entry
1318        while !pager_and_verifier.cache.is_merkle_initialized(key) {
1319            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1320        }
1321
1322        // Resume Open and wait for `create_vmo` to finish
1323        resume_open_tx.send(()).expect("resume_open_tx failed");
1324        let vmo = create_task.await.expect("create_vmo failed");
1325
1326        // We should be able to deliver data and read it immediately
1327        let chunk = &blob_data[..DELIVERY_DATA_SIZE];
1328        let mut data_payload = sender.reserve_payload(chunk.len()).expect("reserve_payload failed");
1329        data_payload.data().copy_from_slice(chunk);
1330        let data_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1331            key: key as u64,
1332            offset: data_payload.offset(),
1333            length: chunk.len() as u32,
1334            target_offset: 0,
1335        }
1336        .into();
1337        data_payload.commit(data_cmd).expect("commit failed");
1338
1339        let mut read_buf = vec![0u8; DELIVERY_DATA_SIZE];
1340        vmo.read(&mut read_buf, 0).expect("read failed");
1341        assert_eq!(read_buf, chunk);
1342
1343        drop(vmo);
1344        close_rx.await.expect("close_rx failed");
1345
1346        drop(pager_and_verifier);
1347        mapping_task.await;
1348        mapper_task.await;
1349    }
1350
1351    #[fuchsia::test]
1352    async fn test_delivery_register_blob_verification() {
1353        let env = TestEnv::new(TEST_BLOB_SIZE).await;
1354
1355        let _paged_vmo =
1356            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1357
1358        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1359            env.delivery_vmo
1360                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1361                .expect("duplicate_handle failed"),
1362            zx::system_get_page_size() as usize, // alignment of payload
1363            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1364        )
1365        .expect("SyncSender::new failed");
1366
1367        // Test corrupted leaves should fail register blob
1368        let mut corrupted_leaves = env.valid_leaves.clone();
1369        corrupted_leaves[0] ^= 0xFF;
1370
1371        let mut bad_payload =
1372            sender.reserve_payload(corrupted_leaves.len()).expect("reserve_payload failed");
1373        bad_payload.data().copy_from_slice(&corrupted_leaves);
1374
1375        let bad_raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1376            key: TEST_VMO_KEY,
1377            offset: bad_payload.offset(),
1378            length: corrupted_leaves.len() as u32,
1379        }
1380        .into();
1381        bad_payload.commit(bad_raw_cmd).expect("commit failed");
1382
1383        // Yield to allow background thread to process the invalid merkle
1384        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1385
1386        let blob =
1387            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1388
1389        // merkle_verifier should remain as None as the leaves are corrupted
1390        assert!(blob.merkle_verifier.get().is_none());
1391
1392        // Test with valid leaves
1393        let mut payload =
1394            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1395        payload.data().copy_from_slice(&env.valid_leaves);
1396
1397        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1398            key: TEST_VMO_KEY,
1399            offset: payload.offset(),
1400            length: env.valid_leaves.len() as u32,
1401        }
1402        .into();
1403        payload.commit(raw_cmd).expect("commit failed");
1404
1405        while blob.merkle_verifier.get().is_none() {
1406            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1407        }
1408
1409        drop(_paged_vmo);
1410        env.teardown().await;
1411    }
1412
1413    #[fuchsia::test]
1414    async fn test_register_blob_invalid_commands() {
1415        let env = TestEnv::new(TEST_BLOB_SIZE).await;
1416
1417        let _paged_vmo =
1418            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1419
1420        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1421            env.delivery_vmo
1422                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1423                .expect("duplicate_handle failed"),
1424            std::mem::align_of::<RawDeliveryCommand>(),
1425            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1426        )
1427        .expect("SyncSender::new failed");
1428
1429        let blob =
1430            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1431
1432        // Test with unknown/expired key
1433        let mut invalid_key_payload =
1434            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1435        invalid_key_payload.data().copy_from_slice(&env.valid_leaves);
1436        let invalid_key_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1437            key: 9999 as u64,
1438            offset: invalid_key_payload.offset(),
1439            length: env.valid_leaves.len() as u32,
1440        }
1441        .into();
1442        invalid_key_payload.commit(invalid_key_cmd).expect("commit failed");
1443        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1444        assert!(blob.merkle_verifier.get().is_none());
1445
1446        // Test out-of-bounds offset/length
1447        let oob_payload = sender.reserve_payload(8).expect("reserve_payload failed");
1448        let oob_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1449            key: TEST_VMO_KEY,
1450            offset: u32::MAX - 4, // Malicious offset
1451            length: 8,
1452        }
1453        .into();
1454        oob_payload.commit(oob_cmd).expect("commit failed");
1455        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1456        assert!(blob.merkle_verifier.get().is_none());
1457
1458        // Test invalid leaf length (not a multiple of HASH_SIZE)
1459        let invalid_len_payload = sender.reserve_payload(10).expect("reserve_payload failed");
1460        let invalid_len_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1461            key: TEST_VMO_KEY,
1462            offset: invalid_len_payload.offset(),
1463            length: 10,
1464        }
1465        .into();
1466        invalid_len_payload.commit(invalid_len_cmd).expect("commit failed");
1467        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1468        assert!(blob.merkle_verifier.get().is_none());
1469
1470        drop(_paged_vmo);
1471        env.teardown().await;
1472    }
1473
1474    #[fuchsia::test]
1475    async fn test_delivery_data_supplies_pages() {
1476        let env = TestEnv::new(TEST_BLOB_SIZE).await;
1477
1478        let paged_vmo =
1479            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1480
1481        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1482            env.delivery_vmo
1483                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1484                .expect("duplicate_handle failed"),
1485            zx::system_get_page_size() as usize, // alignment of payload
1486            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1487        )
1488        .expect("SyncSender::new failed");
1489
1490        let mut payload =
1491            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1492        payload.data().copy_from_slice(&env.valid_leaves);
1493        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1494            key: TEST_VMO_KEY,
1495            offset: payload.offset(),
1496            length: env.valid_leaves.len() as u32,
1497        }
1498        .into();
1499        payload.commit(raw_cmd).expect("commit failed");
1500
1501        let blob =
1502            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1503        while blob.merkle_verifier.get().is_none() {
1504            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1505        }
1506
1507        let test_payload = &env.blob_data[..DELIVERY_DATA_SIZE];
1508
1509        let mut payload =
1510            sender.reserve_payload(test_payload.len()).expect("reserve_payload failed");
1511        payload.data().copy_from_slice(&test_payload);
1512
1513        let cmd: RawDeliveryCommand = DeliveryCommand::Data {
1514            key: TEST_VMO_KEY,
1515            target_offset: 0,
1516            length: test_payload.len() as u32,
1517            offset: payload.offset(),
1518        }
1519        .into();
1520        payload.commit(cmd).expect("commit failed");
1521
1522        let mut read_buf = vec![0u8; DELIVERY_DATA_SIZE];
1523        paged_vmo.read(&mut read_buf, 0).expect("read paged_vmo failed");
1524        assert_eq!(read_buf, test_payload);
1525
1526        drop(paged_vmo);
1527        env.teardown().await;
1528    }
1529
1530    #[fuchsia::test]
1531    async fn test_delivery_data_corrupted() {
1532        // Verify that when a delivered chunk fails verification, client threads waiting
1533        // for data within that chunk receive `ZX_ERR_IO_DATA_INTEGRITY`.
1534        let blob_size = DELIVERY_DATA_SIZE as u64;
1535        let env = TestEnv::new(blob_size).await;
1536
1537        let paged_vmo =
1538            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1539
1540        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1541            env.delivery_vmo
1542                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1543                .expect("duplicate_handle failed"),
1544            zx::system_get_page_size() as usize,
1545            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1546        )
1547        .expect("SyncSender::new failed");
1548
1549        let mut payload =
1550            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1551        payload.data().copy_from_slice(&env.valid_leaves);
1552        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1553            key: TEST_VMO_KEY,
1554            offset: payload.offset(),
1555            length: env.valid_leaves.len() as u32,
1556        }
1557        .into();
1558        payload.commit(raw_cmd).expect("commit failed");
1559
1560        let blob =
1561            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1562        while blob.merkle_verifier.get().is_none() {
1563            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1564        }
1565
1566        // Prepare a corrupted 128 KiB chunk (the entire blob data).
1567        let mut corrupted_data = env.blob_data.clone();
1568        corrupted_data[0] ^= 0xFF;
1569
1570        let mut payload =
1571            sender.reserve_payload(corrupted_data.len()).expect("reserve_payload failed");
1572        payload.data().copy_from_slice(&corrupted_data);
1573        let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1574            key: TEST_VMO_KEY,
1575            offset: payload.offset(),
1576            length: corrupted_data.len() as u32,
1577            target_offset: 0,
1578        }
1579        .into();
1580
1581        // Start client threads reading different pages within the delivery data chunk.
1582        // When verification fails, all waiting threads must fail with `IO_DATA_INTEGRITY`.
1583        let page_size = zx::system_get_page_size() as usize;
1584        let mut reader1 =
1585            spawn_reader_expect_error(&paged_vmo, 0, page_size, zx::Status::IO_DATA_INTEGRITY);
1586        let mut reader2 = spawn_reader_expect_error(
1587            &paged_vmo,
1588            page_size as u64,
1589            page_size,
1590            zx::Status::IO_DATA_INTEGRITY,
1591        );
1592
1593        // Wait for both client threads to start executing before giving them time to block on their
1594        // page faults.
1595        reader1.wait_started().await;
1596        reader2.wait_started().await;
1597        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1598
1599        // Commit the corrupt chunk. Verification will fail, causing the blocked `vmo.read()`
1600        // calls for that chunk to return `IO_DATA_INTEGRITY`.
1601        payload.commit(raw_cmd).expect("commit failed");
1602
1603        reader1.wait_and_verify().await;
1604        reader2.wait_and_verify().await;
1605
1606        drop(paged_vmo);
1607        env.teardown().await;
1608    }
1609
1610    #[fuchsia::test]
1611    async fn test_delivery_data_corrupted_chunk_preserves_supplied_pages() {
1612        let chunk_size = DELIVERY_DATA_SIZE;
1613        let blob_size = (chunk_size * 2) as u64;
1614        let env = TestEnv::new(blob_size).await;
1615
1616        let paged_vmo =
1617            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1618
1619        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1620            env.delivery_vmo
1621                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1622                .expect("duplicate_handle failed"),
1623            zx::system_get_page_size() as usize,
1624            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1625        )
1626        .expect("SyncSender::new failed");
1627
1628        let mut payload =
1629            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1630        payload.data().copy_from_slice(&env.valid_leaves);
1631        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1632            key: TEST_VMO_KEY,
1633            offset: payload.offset(),
1634            length: env.valid_leaves.len() as u32,
1635        }
1636        .into();
1637        payload.commit(raw_cmd).expect("commit failed");
1638
1639        let blob =
1640            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1641        while blob.merkle_verifier.get().is_none() {
1642            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1643        }
1644
1645        let page_size = zx::system_get_page_size() as usize;
1646
1647        // Chunk 1 test for successful page request.
1648        // Start a reader on page 0 and deliver the valid first chunk (0..chunk_size).
1649        let paged_vmo_clone1 =
1650            paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate_handle failed");
1651        let (tx1, rx1) = oneshot::channel();
1652        let expected_page0 = env.blob_data[..page_size].to_vec();
1653        let thread1 = std::thread::spawn(move || {
1654            let mut buf = vec![0u8; page_size];
1655            paged_vmo_clone1.read(&mut buf, 0).expect("read should succeed");
1656            assert_eq!(buf, expected_page0);
1657            let _ = tx1.send(());
1658        });
1659        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1660
1661        let valid_chunk = &env.blob_data[..chunk_size];
1662        let mut payload =
1663            sender.reserve_payload(valid_chunk.len()).expect("reserve_payload failed");
1664        payload.data().copy_from_slice(valid_chunk);
1665        let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1666            key: TEST_VMO_KEY,
1667            offset: payload.offset(),
1668            length: valid_chunk.len() as u32,
1669            target_offset: 0,
1670        }
1671        .into();
1672        payload.commit(raw_cmd).expect("commit failed");
1673
1674        rx1.await.expect("thread1 panicked or hung");
1675        thread1.join().expect("join failed");
1676
1677        // Test for failed page requests on the subsequent chunk.
1678        let mut reader2 = spawn_reader_expect_error(
1679            &paged_vmo,
1680            chunk_size as u64,
1681            page_size,
1682            zx::Status::IO_DATA_INTEGRITY,
1683        );
1684        reader2.wait_started().await;
1685        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1686
1687        // Deliver corrupted second chunk (chunk_size..chunk_size * 2).
1688        let mut corrupted_chunk = env.blob_data[chunk_size..].to_vec();
1689        corrupted_chunk[0] ^= 0xFF;
1690
1691        let mut payload =
1692            sender.reserve_payload(corrupted_chunk.len()).expect("reserve_payload failed");
1693        payload.data().copy_from_slice(&corrupted_chunk);
1694        let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1695            key: TEST_VMO_KEY,
1696            offset: payload.offset(),
1697            length: corrupted_chunk.len() as u32,
1698            target_offset: chunk_size as u64,
1699        }
1700        .into();
1701        payload.commit(raw_cmd).expect("commit failed");
1702
1703        reader2.wait_and_verify().await;
1704
1705        // Verify that the previously supplied page 0 remains intact and readable.
1706        let mut buf = vec![0u8; page_size];
1707        paged_vmo.read(&mut buf, 0).expect("previously supplied page 0 must still be readable");
1708        assert_eq!(buf, env.blob_data[..page_size]);
1709
1710        drop(paged_vmo);
1711        env.teardown().await;
1712    }
1713
1714    #[fuchsia::test]
1715    async fn test_delivery_data_corrupted_multi_chunk_read() {
1716        let chunk_size = DELIVERY_DATA_SIZE;
1717        let blob_size = (chunk_size * 3) as u64;
1718        let env = TestEnv::new(blob_size).await;
1719
1720        let paged_vmo =
1721            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1722
1723        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1724            env.delivery_vmo
1725                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1726                .expect("duplicate_handle failed"),
1727            zx::system_get_page_size() as usize,
1728            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1729        )
1730        .expect("SyncSender::new failed");
1731
1732        let mut payload =
1733            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1734        payload.data().copy_from_slice(&env.valid_leaves);
1735        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1736            key: TEST_VMO_KEY,
1737            offset: payload.offset(),
1738            length: env.valid_leaves.len() as u32,
1739        }
1740        .into();
1741        payload.commit(raw_cmd).expect("commit failed");
1742
1743        let blob =
1744            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1745        while blob.merkle_verifier.get().is_none() {
1746            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1747        }
1748
1749        // A 384 KiB `vmo.read()` triggers a single page request spanning multiple pages.
1750        // Chunk 1 (0..128 KiB) is valid, Chunk 2 (128..256 KiB) is corrupted, and Chunk 3 is not
1751        // delivered because delivery stops on verification failure.
1752        // Verify that when Chunk 2 fails verification, the page request fails and
1753        // `vmo.read()` returns `IO_DATA_INTEGRITY`.
1754        let mut reader = spawn_reader_expect_error(
1755            &paged_vmo,
1756            0,
1757            blob_size as usize,
1758            zx::Status::IO_DATA_INTEGRITY,
1759        );
1760        reader.wait_started().await;
1761        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1762
1763        // Deliver Chunk 1 (0..128 KiB) valid.
1764        let valid_chunk1 = &env.blob_data[..chunk_size];
1765        let mut payload =
1766            sender.reserve_payload(valid_chunk1.len()).expect("reserve_payload failed");
1767        payload.data().copy_from_slice(valid_chunk1);
1768        let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1769            key: TEST_VMO_KEY,
1770            offset: payload.offset(),
1771            length: valid_chunk1.len() as u32,
1772            target_offset: 0,
1773        }
1774        .into();
1775        payload.commit(raw_cmd).expect("commit failed");
1776
1777        // Allow the reader thread time to receive the supplied first chunk, copy it, and block on
1778        // the second chunk's unpopulated pages.
1779        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1780
1781        // Deliver Chunk 2 (128..256 KiB) corrupted.
1782        let mut corrupted_chunk2 = env.blob_data[chunk_size..chunk_size * 2].to_vec();
1783        corrupted_chunk2[0] ^= 0xFF;
1784        let mut payload =
1785            sender.reserve_payload(corrupted_chunk2.len()).expect("reserve_payload failed");
1786        payload.data().copy_from_slice(&corrupted_chunk2);
1787        let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1788            key: TEST_VMO_KEY,
1789            offset: payload.offset(),
1790            length: corrupted_chunk2.len() as u32,
1791            target_offset: chunk_size as u64,
1792        }
1793        .into();
1794        payload.commit(raw_cmd).expect("commit failed");
1795
1796        // Don't deliver chunk 3 since chunk 2 failed.
1797
1798        // The `vmo.read()` must fail with IO_DATA_INTEGRITY.
1799        reader.wait_and_verify().await;
1800
1801        drop(paged_vmo);
1802        env.teardown().await;
1803    }
1804
1805    #[fuchsia::test]
1806    async fn test_delivery_data_uninitialized() {
1807        let env = TestEnv::new(TEST_BLOB_SIZE).await;
1808
1809        let paged_vmo =
1810            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1811
1812        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1813            env.delivery_vmo
1814                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1815                .expect("duplicate_handle failed"),
1816            zx::system_get_page_size() as usize,
1817            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1818        )
1819        .expect("SyncSender::new failed");
1820
1821        // Spawn reader thread to generate the initial page request
1822        let mut reader = spawn_reader_expect_error(&paged_vmo, 0, 8192, zx::Status::BAD_STATE);
1823        reader.wait_started().await;
1824        fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1825
1826        // Intentionally push a Data chunk without sending RegisterBlob first
1827        let chunk = &env.blob_data[..8192];
1828        let mut payload = sender.reserve_payload(chunk.len()).expect("reserve_payload failed");
1829        payload.data().copy_from_slice(chunk);
1830        let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1831            key: TEST_VMO_KEY,
1832            offset: payload.offset(),
1833            length: chunk.len() as u32,
1834            target_offset: 0,
1835        }
1836        .into();
1837        payload.commit(raw_cmd).expect("commit failed");
1838
1839        reader.wait_and_verify().await;
1840
1841        drop(paged_vmo);
1842        env.teardown().await;
1843    }
1844
1845    #[fuchsia::test]
1846    async fn test_delivery_data_multiple_chunks() {
1847        let env = TestEnv::new(TEST_BLOB_SIZE).await;
1848
1849        let paged_vmo =
1850            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1851
1852        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1853            env.delivery_vmo
1854                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1855                .expect("duplicate_handle failed"),
1856            zx::system_get_page_size() as usize,
1857            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1858        )
1859        .expect("SyncSender::new failed");
1860
1861        let mut payload =
1862            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1863        payload.data().copy_from_slice(&env.valid_leaves);
1864        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1865            key: TEST_VMO_KEY,
1866            offset: payload.offset(),
1867            length: env.valid_leaves.len() as u32,
1868        }
1869        .into();
1870        payload.commit(raw_cmd).expect("commit failed");
1871
1872        let blob =
1873            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1874        while blob.merkle_verifier.get().is_none() {
1875            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1876        }
1877
1878        let expected_data = env.blob_data.clone();
1879        let paged_vmo_clone =
1880            paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate_handle failed");
1881
1882        let (tx, rx) = oneshot::channel();
1883        let thread = std::thread::spawn(move || {
1884            let mut buf = vec![0u8; expected_data.len()];
1885            paged_vmo_clone.read(&mut buf, 0).expect("failed to read from paged vmo");
1886            assert_eq!(buf, expected_data);
1887            let _ = tx.send(());
1888        });
1889
1890        // Push incrementally - chunks must be a multiple of DELIVERY_DATA_SIZE
1891        let chunk_size = DELIVERY_DATA_SIZE;
1892        for (i, chunk) in env.blob_data.chunks(chunk_size).enumerate() {
1893            let mut payload = sender.reserve_payload(chunk.len()).expect("reserve_payload failed");
1894            payload.data().copy_from_slice(chunk);
1895            let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1896                key: TEST_VMO_KEY,
1897                offset: payload.offset(),
1898                length: chunk.len() as u32,
1899                target_offset: (i * chunk_size) as u64,
1900            }
1901            .into();
1902            payload.commit(raw_cmd).expect("commit failed");
1903        }
1904
1905        rx.await.expect("reading thread panicked or hung");
1906        thread.join().expect("join failed");
1907
1908        drop(paged_vmo);
1909        env.teardown().await;
1910    }
1911
1912    #[fuchsia::test]
1913    async fn test_delivery_data_unaligned_blob_size() {
1914        let data_size = (DELIVERY_DATA_SIZE + 1024) as u64; // Blob with unaligned size
1915        let env = TestEnv::new(data_size).await;
1916
1917        let paged_vmo =
1918            env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("Failed to create VMO");
1919
1920        let mut sender = SyncSender::<RawDeliveryCommand>::new(
1921            env.delivery_vmo
1922                .duplicate_handle(zx::Rights::SAME_RIGHTS)
1923                .expect("duplicate_handle failed"),
1924            zx::system_get_page_size() as usize, // alignment of payload
1925            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1926        )
1927        .expect("SyncSender::new failed");
1928
1929        let mut payload =
1930            sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1931        payload.data().copy_from_slice(&env.valid_leaves);
1932        let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1933            key: TEST_VMO_KEY,
1934            offset: payload.offset(),
1935            length: env.valid_leaves.len() as u32,
1936        }
1937        .into();
1938        payload.commit(raw_cmd).expect("commit failed");
1939
1940        let blob =
1941            env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1942        while blob.merkle_verifier.get().is_none() {
1943            fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1944        }
1945
1946        let test_payload = env.blob_data.clone();
1947        assert_eq!(test_payload.len(), data_size as usize);
1948
1949        let chunk1 = &test_payload[..DELIVERY_DATA_SIZE];
1950        let mut payload1 = sender.reserve_payload(chunk1.len()).expect("reserve_payload failed");
1951        payload1.data().copy_from_slice(chunk1);
1952        let off1 = payload1.offset();
1953        payload1
1954            .commit(
1955                DeliveryCommand::Data {
1956                    key: TEST_VMO_KEY,
1957                    target_offset: 0,
1958                    length: chunk1.len() as u32,
1959                    offset: off1,
1960                }
1961                .into(),
1962            )
1963            .expect("commit failed");
1964
1965        let chunk2 = &test_payload[DELIVERY_DATA_SIZE..];
1966        let page_size = zx::system_get_page_size() as usize;
1967        let chunk2_len_aligned = chunk2.len().div_ceil(page_size) * page_size;
1968        let mut payload2 =
1969            sender.reserve_payload(chunk2_len_aligned).expect("reserve_payload failed");
1970        let payload2_data = payload2.data();
1971        payload2_data.subslice_mut(0..chunk2.len()).copy_from_slice(chunk2);
1972        payload2_data.subslice_mut(chunk2.len()..chunk2_len_aligned).fill(0);
1973
1974        let off2 = payload2.offset();
1975        payload2
1976            .commit(
1977                DeliveryCommand::Data {
1978                    key: TEST_VMO_KEY,
1979                    target_offset: DELIVERY_DATA_SIZE as u64,
1980                    length: chunk2_len_aligned as u32,
1981                    offset: off2,
1982                }
1983                .into(),
1984            )
1985            .expect("commit failed");
1986
1987        let mut read_buf = vec![0u8; data_size as usize];
1988        paged_vmo.read(&mut read_buf, 0).expect("read paged_vmo failed");
1989        assert_eq!(read_buf, test_payload);
1990
1991        let page_size = zx::system_get_page_size() as u64;
1992        // Round the unaligned data size to the next page boundary
1993        let out_of_bounds_offset = data_size.div_ceil(page_size) * page_size;
1994        // Attempt to read from an out of bounds offset
1995        let mut read_buf = vec![0u8; 1];
1996        let res = paged_vmo.read(&mut read_buf, out_of_bounds_offset);
1997        assert_eq!(res, Err(zx::Status::OUT_OF_RANGE));
1998
1999        drop(paged_vmo);
2000        env.teardown().await;
2001    }
2002}