Skip to main content

fxfs_platform/fuchsia/fxblob/
mapping_provider.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::fuchsia::errors::map_to_status;
6use crate::fuchsia::fxblob::blob::FxBlob;
7use crate::fuchsia::fxblob::directory::BlobDirectory;
8use crate::fuchsia::node::OpenedNode;
9use crate::fuchsia::pager::PagerBacked;
10use anyhow::Error;
11use fidl_fuchsia_storage_mapping as fmapping;
12use fuchsia_merkle::Hash;
13use futures::TryStreamExt;
14use fxfs::errors::FxfsError;
15use log::{error, warn};
16use mapping::{
17    Extents, MAPPING_VMO_SIZE, MappingCommand, PENDING_COMMANDS_CAPACITY, RawMappingCommand,
18};
19use std::collections::HashMap;
20use std::sync::Arc;
21use vmo_fifo::AsyncSender;
22
23/// BlobMappingProvider services requests to open mapping sessions over a given BlobDirectory.
24pub struct BlobMappingProvider {
25    blob_directory: Arc<BlobDirectory>,
26}
27
28impl BlobMappingProvider {
29    pub fn new(blob_directory: Arc<BlobDirectory>) -> Result<Self, Error> {
30        Ok(Self { blob_directory })
31    }
32
33    pub async fn handle_mapping_provider_requests(
34        self: Arc<Self>,
35        mut stream: fmapping::MappingProviderRequestStream,
36    ) {
37        while let Ok(Some(request)) = stream.try_next().await {
38            match request {
39                fmapping::MappingProviderRequest::OpenSession { session, responder } => {
40                    let vmo = match zx::Vmo::create(MAPPING_VMO_SIZE) {
41                        Ok(v) => v,
42                        Err(e) => {
43                            let _ = responder.send(Err(e.into_raw()));
44                            continue;
45                        }
46                    };
47                    let sender = match AsyncSender::<RawMappingCommand>::new(
48                        vmo,
49                        8,
50                        PENDING_COMMANDS_CAPACITY,
51                    ) {
52                        Ok(s) => s,
53                        Err(e) => {
54                            let _ = responder.send(Err(e.into_raw()));
55                            continue;
56                        }
57                    };
58                    let vmo_clone = match sender.vmo().duplicate_handle(zx::Rights::SAME_RIGHTS) {
59                        Ok(v) => v,
60                        Err(e) => {
61                            let _ = responder.send(Err(e.into_raw()));
62                            continue;
63                        }
64                    };
65                    if let Err(error) = responder.send(Ok(vmo_clone)) {
66                        error!(error:?; "Failed to send open session response");
67                    } else {
68                        let mapping_session =
69                            BlobMappingSession::new(self.blob_directory.clone(), sender);
70                        self.blob_directory.volume().scope().spawn(async move {
71                            mapping_session
72                                .handle_mapping_session_requests(session.into_stream())
73                                .await;
74                        });
75                    }
76                }
77                fmapping::MappingProviderRequest::_UnknownMethod { ordinal, .. } => {
78                    warn!(ordinal; "Unknown MappingProvider method");
79                }
80            }
81        }
82    }
83}
84
85/// `BlobMappingSession` services requests to open or close blobs within a `BlobDirectory`.
86/// Upon opening a blob, the blob's underlying extents (data and merkle blocks) are retrieved and
87/// then forwarded to client asynchronously via `sender`. The client provides a session-unique
88/// `key` for each opened blob, which can also be referenced during a `Close` request.
89pub struct BlobMappingSession {
90    blob_directory: Arc<BlobDirectory>,
91    sender: AsyncSender<RawMappingCommand>,
92    opened_blobs: HashMap<u64, OpenedNode<FxBlob>>,
93}
94
95impl BlobMappingSession {
96    pub fn new(blob_directory: Arc<BlobDirectory>, sender: AsyncSender<RawMappingCommand>) -> Self {
97        Self { blob_directory, sender, opened_blobs: HashMap::new() }
98    }
99
100    /// Retrieves the extent mappings for the blob and registers the blob in the mapping session.
101    /// Returns the uncompressed byte size of the blob.
102    async fn open_blob(&mut self, key: u64, hash: Hash) -> Result<u64, Error> {
103        let node = self.blob_directory.open_blob(&hash.into()).await?.ok_or(FxfsError::NotFound)?;
104
105        let extents = node.get_mapping_extents().await?;
106        let size = node.as_ref().byte_size();
107        let stored_size = node.as_ref().stored_size().await?;
108        let extent_count = extents.data.len() as u32;
109        let metadata_count = extents.merkle.len() as u32;
110
111        let allocation_size = (extent_count + metadata_count) as usize * std::mem::size_of::<u64>();
112
113        if allocation_size > 0 {
114            let mut payload = self.sender.reserve_payload(allocation_size).await?;
115            let offset_in_vmo = payload.offset();
116
117            let data_extents = Extents::try_new(&extents.data, 0)?;
118            let merkle_extents = Extents::try_new(&extents.merkle, 0)?;
119
120            for (mut chunk, val_res) in payload.data().chunks_mut(std::mem::size_of::<u64>()).zip(
121                Extents::encode_extents(&data_extents)
122                    .chain(Extents::encode_extents(&merkle_extents)),
123            ) {
124                chunk.copy_from_slice(&val_res.to_le_bytes());
125            }
126
127            let command = MappingCommand::Mappings {
128                key,
129                offset: offset_in_vmo as u32,
130                stored_size,
131                device_offset: 0,
132                metadata_count,
133                extent_count,
134                encrypted: false,
135            };
136
137            payload.commit(command.into()).await?;
138        }
139
140        self.opened_blobs.insert(key, node);
141
142        Ok(size)
143    }
144
145    /// Unregisters the blob mapping and signals the block driver to terminate tracking.
146    async fn close_blob(&mut self, key: u64) -> Result<(), Error> {
147        if self.opened_blobs.remove(&key).is_none() {
148            return Ok(());
149        }
150        self.sender.push(MappingCommand::CloseBlob { key }.into()).await?;
151        Ok(())
152    }
153
154    pub async fn handle_mapping_session_requests(
155        mut self,
156        mut stream: fmapping::MappingSessionRequestStream,
157    ) {
158        while let Ok(Some(request)) = stream.try_next().await {
159            match request {
160                fmapping::MappingSessionRequest::Open { key, identifier, responder } => {
161                    // We expect the identifier to be a Merkle Root Hash for the blob.
162                    if identifier.len() != fuchsia_hash::HASH_SIZE {
163                        responder.send(Err(zx::Status::INVALID_ARGS.into_raw())).unwrap_or_else(
164                            |error| warn!(error:?; "Failed to send mapping session response"),
165                        );
166                        continue;
167                    }
168                    let hash = Hash::from(<[u8; 32]>::try_from(identifier.as_slice()).unwrap());
169                    match self.open_blob(key, hash).await {
170                        Ok(size) => {
171                            responder.send(Ok(size)).unwrap_or_else(
172                                |error| warn!(error:?; "Failed to send mapping session response"),
173                            );
174                        }
175                        Err(error) => {
176                            error!(error:?; "Failed to open blob");
177                            responder.send(Err(map_to_status(error).into_raw())).unwrap_or_else(
178                                |error| warn!(error:?; "Failed to send mapping session response"),
179                            );
180                        }
181                    }
182                }
183                fmapping::MappingSessionRequest::Close { key, responder } => {
184                    let result = match self.close_blob(key).await {
185                        Ok(()) => Ok(()),
186                        Err(error) => {
187                            error!(error:?; "Failed to close blob");
188                            Err(map_to_status(error).into_raw())
189                        }
190                    };
191                    responder.send(result).unwrap_or_else(
192                        |error| warn!(error:?; "Failed to send mapping session response"),
193                    );
194                }
195                fmapping::MappingSessionRequest::_UnknownMethod { ordinal, .. } => {
196                    warn!(ordinal; "Unknown MappingSession method");
197                }
198            }
199        }
200    }
201}
202
203#[cfg(test)]
204mod tests {
205    use super::*;
206    use crate::fuchsia::fxblob::testing::{BlobFixture, new_blob_fixture, open_blob_fixture};
207    use crate::fuchsia::testing::TestFixture;
208    use blob_writer::BlobWriter;
209    use delivery_blob::{CompressionMode, Type1Blob};
210    use fidl_fuchsia_io::UnlinkOptions;
211    use fuchsia_async as fasync;
212    use futures::channel::oneshot;
213    use fxfs::object_handle::ObjectHandle;
214    use fxfs::object_store::{HandleOptions, ObjectStore};
215    use storage_device::Device;
216    use storage_device::buffer::OwnedBuffer;
217    use storage_device::buffer_allocator::{BufferAllocator, BufferSource};
218    use vmo_fifo::Receiver;
219
220    #[fuchsia::test]
221    async fn test_blob_mapping_provider() {
222        let fixture = new_blob_fixture().await;
223        // Test with a large amount of non-compressible data to generate many extents
224        let data = vec![42; 300_000];
225        let hash = fixture.write_blob(&data, CompressionMode::Never).await;
226
227        let blob_dir = fixture
228            .volume()
229            .root()
230            .clone()
231            .as_node()
232            .into_any()
233            .downcast::<BlobDirectory>()
234            .expect("Failed to downcast root directory to BlobDirectory");
235
236        let node = blob_dir
237            .open_blob(&hash.into())
238            .await
239            .expect("Failed to open blob in Fxfs")
240            .expect("open_blob returned None instead of node");
241        let extents = node.get_mapping_extents().await.expect("Failed to retrieve extents");
242        let data_extents = extents.data;
243        let merkle_extents = extents.merkle;
244
245        // Un-dropped nodes pin the Blob as actively opened. The unmount routine in fixture.close()
246        // will wait forever for this blob to be fully closed, causing a test timeout.
247        drop(node);
248
249        let vmo = zx::Vmo::create(MAPPING_VMO_SIZE).unwrap();
250        let client_mapping = vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap();
251        let sender =
252            AsyncSender::<RawMappingCommand>::new(vmo, 8, PENDING_COMMANDS_CAPACITY).unwrap();
253
254        let mut session = BlobMappingSession::new(blob_dir, sender);
255
256        let receiver_task = fasync::unblock(move || {
257            let mut receiver = Receiver::<RawMappingCommand>::new(client_mapping, 256)
258                .expect("Failed to create the Receiver wrapper");
259
260            // First Open Command
261            let cmd1_raw = receiver.peek().expect("peek failed");
262            let cmd1 = MappingCommand::try_from(*cmd1_raw).expect("try_from failed");
263
264            let (cmd1_offset, cmd1_extent_count, cmd1_metadata_count) = match cmd1 {
265                MappingCommand::Mappings {
266                    key,
267                    offset,
268                    stored_size: _,
269                    device_offset: _,
270                    metadata_count,
271                    extent_count,
272                    encrypted,
273                } => {
274                    assert_eq!(key, 1);
275                    assert_eq!(extent_count, data_extents.len() as u32);
276                    assert_eq!(metadata_count, merkle_extents.len() as u32);
277                    assert!(!encrypted);
278                    (offset, extent_count, metadata_count)
279                }
280                _ => panic!("Expected Mappings command"),
281            };
282
283            // Verify payload
284            let total_extents = cmd1_extent_count + cmd1_metadata_count;
285            let buffer = cmd1_raw.payload_slice(cmd1_offset, total_extents * 8).to_vec();
286
287            let data_extents_container = Extents::try_new(&data_extents, 0).unwrap();
288            let merkle_extents_container = Extents::try_new(&merkle_extents, 0).unwrap();
289            let mut expected_payload = Vec::new();
290            for val in Extents::encode_extents(&data_extents_container)
291                .chain(Extents::encode_extents(&merkle_extents_container))
292            {
293                expected_payload.extend_from_slice(&val.to_le_bytes());
294            }
295            assert_eq!(buffer, expected_payload);
296            cmd1_raw.pop().expect("Failed pop_commit");
297
298            let cmd2_raw = receiver.peek().expect("peek failed");
299            let cmd2 = MappingCommand::try_from(*cmd2_raw).expect("try_from failed");
300            match cmd2 {
301                MappingCommand::CloseBlob { key } => assert_eq!(key, 1),
302                _ => panic!("Expected CloseBlob command"),
303            };
304            cmd2_raw.pop().expect("Failed pop_commit");
305        });
306
307        let server_task = async move {
308            let key = 1;
309            let size = session.open_blob(key, hash).await.expect("open_blob failed");
310            assert_eq!(size, data.len() as u64);
311
312            session.close_blob(key).await.expect("close_blob failed on existing key");
313
314            std::mem::drop(session);
315        };
316
317        futures::join!(receiver_task, server_task);
318
319        fixture.close().await;
320    }
321
322    #[fuchsia::test]
323    async fn test_missing_blob() {
324        let fixture = new_blob_fixture().await;
325        let blob_dir = fixture
326            .volume()
327            .root()
328            .clone()
329            .as_node()
330            .into_any()
331            .downcast::<BlobDirectory>()
332            .expect("Failed to downcast");
333
334        let vmo = zx::Vmo::create(MAPPING_VMO_SIZE).unwrap();
335        let sender =
336            AsyncSender::<RawMappingCommand>::new(vmo, 8, PENDING_COMMANDS_CAPACITY).unwrap();
337        let mut session = BlobMappingSession::new(blob_dir, sender);
338        let hash = Hash::from([1u8; 32]);
339        session
340            .open_blob(1, hash)
341            .await
342            .expect_err("open_blob should fail with blob that doesn't exist");
343
344        std::mem::drop(session);
345        fixture.close().await;
346    }
347
348    #[fuchsia::test]
349    async fn test_invalid_key() {
350        let fixture = new_blob_fixture().await;
351        let blob_dir = fixture
352            .volume()
353            .root()
354            .clone()
355            .as_node()
356            .into_any()
357            .downcast::<BlobDirectory>()
358            .expect("Failed to downcast");
359
360        let vmo = zx::Vmo::create(MAPPING_VMO_SIZE).unwrap();
361        let sender =
362            AsyncSender::<RawMappingCommand>::new(vmo, 8, PENDING_COMMANDS_CAPACITY).unwrap();
363        let mut session = BlobMappingSession::new(blob_dir, sender);
364        session.close_blob(42).await.expect("close_blob should return Ok with invalid key");
365        std::mem::drop(session);
366
367        fixture.close().await;
368    }
369
370    // Verifies that unlinking an open blob does not deallocate its extents until all references
371    // are closed.
372    #[fuchsia::test]
373    async fn test_unlink_open_blob() {
374        let fixture = new_blob_fixture().await;
375        let data = vec![77u8; 8192];
376        let hash = fixture.write_blob(&data, CompressionMode::Never).await;
377        let hash_string = format!("{}", hash);
378
379        let object_id = fixture.get_blob_handle(&hash_string).await.object_id();
380
381        let blob_dir = fixture
382            .volume()
383            .root()
384            .clone()
385            .as_node()
386            .into_any()
387            .downcast::<BlobDirectory>()
388            .expect("Failed to downcast root directory to BlobDirectory");
389
390        let vmo = zx::Vmo::create(MAPPING_VMO_SIZE).unwrap();
391        let sender =
392            AsyncSender::<RawMappingCommand>::new(vmo, 8, PENDING_COMMANDS_CAPACITY).unwrap();
393        let mut session = BlobMappingSession::new(blob_dir, sender);
394
395        let key = 1;
396        session.open_blob(key, hash).await.expect("open_blob failed");
397
398        // Unlink the blob from the root directory.
399        fixture
400            .root()
401            .unlink(&hash_string, &UnlinkOptions::default())
402            .await
403            .expect("FIDL failed")
404            .expect("unlink failed");
405
406        // Flush the graveyard to trigger purging if the object has no open references.
407        fixture.fs().graveyard().flush().await;
408
409        // Since the blob is currently open in the mapping session, it must not be purged.
410        assert!(
411            ObjectStore::open_object(
412                fixture.volume().volume(),
413                object_id,
414                HandleOptions::default(),
415                None
416            )
417            .await
418            .is_ok(),
419            "Blob was purged from store while still open in BlobMappingSession"
420        );
421
422        // Close the blob in the mapping session.
423        session.close_blob(key).await.expect("close_blob failed");
424
425        // Flush the graveyard again now that the session has released the blob.
426        fixture.fs().graveyard().flush().await;
427
428        // Now the blob must be tombstoned and purged from the object store.
429        assert!(
430            ObjectStore::open_object(
431                fixture.volume().volume(),
432                object_id,
433                HandleOptions::default(),
434                None
435            )
436            .await
437            .is_err(),
438            "Blob should have been purged from store after close_blob"
439        );
440
441        drop(session);
442        fixture.close().await;
443    }
444
445    #[fuchsia::test]
446    async fn test_mapping_provider_and_mapping_session() {
447        let fixture = new_blob_fixture().await;
448        let data = vec![42; 8192];
449        let hash = fixture.write_blob(&data, CompressionMode::Never).await;
450
451        let blob_dir = fixture
452            .volume()
453            .root()
454            .clone()
455            .as_node()
456            .into_any()
457            .downcast::<BlobDirectory>()
458            .expect("Failed to downcast root directory to BlobDirectory");
459
460        let scope = blob_dir.volume().scope().clone();
461        let server = Arc::new(
462            BlobMappingProvider::new(blob_dir).expect("Failed to create BlobMappingProvider"),
463        );
464
465        // Spawn the mapping provider stream.
466        let (provider_proxy, provider_server_end) =
467            fidl::endpoints::create_proxy::<fmapping::MappingProviderMarker>();
468        scope.spawn(async move {
469            server.handle_mapping_provider_requests(provider_server_end.into_stream()).await;
470        });
471
472        // Open a mapping session from the mapping provider
473        let (session_proxy, session_server_end) =
474            fidl::endpoints::create_proxy::<fmapping::MappingSessionMarker>();
475        let _shared_vmo = provider_proxy
476            .open_session(session_server_end)
477            .await
478            .expect("open_session failed")
479            .expect("vmo returned an error");
480
481        let id: [u8; 32] = hash.into();
482        let key = 1;
483        let size = session_proxy
484            .open(key, &id)
485            .await
486            .expect("open failed")
487            .expect("open explicitly returned an error");
488
489        assert_eq!(size, 8192);
490
491        session_proxy
492            .close(key)
493            .await
494            .expect("close failed")
495            .expect("close explicitly returned an error");
496
497        // Test some failures
498
499        // Sending an invalid length hash to open should return INVALID_ARGS
500        let bad_length_id = vec![1u8, 2, 3];
501        let invalid_args_err = session_proxy
502            .open(key, &bad_length_id)
503            .await
504            .expect("open wire call failed")
505            .unwrap_err();
506        assert_eq!(invalid_args_err, zx::Status::INVALID_ARGS.into_raw());
507
508        // Sending a valid length hash that does not exist should return NOT_FOUND
509        let not_found_id = [0u8; 32];
510        let not_found_err = session_proxy
511            .open(key, &not_found_id)
512            .await
513            .expect("open wire call failed")
514            .unwrap_err();
515        assert_eq!(not_found_err, zx::Status::NOT_FOUND.into_raw());
516
517        // Trying to close an invalid blob key currently always succeeds. The BlobMappingServer
518        // unconditionally passes the command down the FIFO queue to the block driver and doesn't
519        // explicitly track active connections.
520        session_proxy.close(123).await.expect("close wire call failed").expect("close failed");
521
522        fixture.close().await;
523    }
524
525    struct DeviceBlockService {
526        device: Arc<dyn Device>,
527    }
528
529    impl mapping::reader::BlockService for DeviceBlockService {
530        fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
531            let block_size = self.device.block_size() as usize;
532            let nblocks = (std::cmp::max(max_len, block_size) + block_size - 1) / block_size;
533            let aligned_len = nblocks * block_size;
534            let pool_size = nblocks.next_power_of_two() * block_size;
535            let buffer_source = BufferSource::new(pool_size);
536            let allocator = BufferAllocator::new(block_size, buffer_source);
537            Arc::new(allocator).allocate_buffer_sync_owned(aligned_len)
538        }
539
540        fn read_blocks(
541            &self,
542            device_offset: u64,
543            mut dest_buffer: OwnedBuffer,
544            on_complete: Box<dyn FnOnce(Result<OwnedBuffer, anyhow::Error>) + Send>,
545        ) -> Result<(), anyhow::Error> {
546            let device = self.device.clone();
547            futures::executor::block_on(async move {
548                let res = device.read(device_offset, dest_buffer.as_mut()).await;
549                on_complete(res.map(|_| dest_buffer));
550            });
551            Ok(())
552        }
553    }
554
555    async fn run_blob_mapping_test(
556        fixture: &TestFixture,
557        test_data: &[u8],
558        mode: CompressionMode,
559        min_data_extents: usize,
560        min_merkle_extents: usize,
561    ) {
562        let hash = fixture.write_blob(test_data, mode).await;
563        run_blob_mapping_test_with_hash(
564            fixture,
565            hash,
566            test_data,
567            min_data_extents,
568            min_merkle_extents,
569        )
570        .await;
571    }
572
573    async fn write_blob_chunked(fx: &TestFixture, data: &[u8], mode: CompressionMode) -> Hash {
574        let hash = fuchsia_merkle::root_from_slice(data);
575        let compressed_data = Type1Blob::generate(data, mode);
576        let writer = fx.create_blob(&hash.into(), false).await.expect("create blob failed");
577        let vmo = writer
578            .get_vmo(compressed_data.len() as u64)
579            .await
580            .expect("transport error on get_vmo")
581            .expect("failed to get vmo");
582        let vmo_size = vmo.get_size().expect("failed to get vmo size");
583
584        // Write in small chunks (512 B) to allocate extents across transactions.
585        let chunk_size = 512;
586        let mut write_offset = 0u64;
587        let mut bytes_left = compressed_data.len() as u64;
588        while bytes_left > 0 {
589            let chunk_len = std::cmp::min(bytes_left, chunk_size);
590            vmo.write(
591                &compressed_data[write_offset as usize..(write_offset + chunk_len) as usize],
592                write_offset % vmo_size,
593            )
594            .expect("failed to write to vmo");
595            let _ = writer
596                .bytes_ready(chunk_len)
597                .await
598                .expect("transport error on bytes_ready")
599                .expect("failed to write data to vmo");
600            write_offset += chunk_len;
601            bytes_left -= chunk_len;
602        }
603        hash
604    }
605
606    async fn run_blob_mapping_test_with_hash(
607        fixture: &TestFixture,
608        hash: Hash,
609        uncompressed_data: &[u8],
610        min_data_extents: usize,
611        min_merkle_extents: usize,
612    ) {
613        let blob_dir = fixture
614            .volume()
615            .root()
616            .clone()
617            .as_node()
618            .into_any()
619            .downcast::<BlobDirectory>()
620            .expect("Failed to downcast root directory to BlobDirectory");
621
622        let scope = blob_dir.volume().scope().clone();
623        let server = Arc::new(
624            BlobMappingProvider::new(blob_dir).expect("Failed to create BlobMappingProvider"),
625        );
626
627        let (provider_proxy, provider_server_end) =
628            fidl::endpoints::create_proxy::<fmapping::MappingProviderMarker>();
629        scope.spawn(async move {
630            server.handle_mapping_provider_requests(provider_server_end.into_stream()).await;
631        });
632
633        let (session_proxy, session_server_end) =
634            fidl::endpoints::create_proxy::<fmapping::MappingSessionMarker>();
635        let mapping_vmo = provider_proxy
636            .open_session(session_server_end)
637            .await
638            .expect("open_session failed")
639            .expect("vmo returned an error");
640
641        let mut receiver =
642            vmo_fifo::Receiver::<RawMappingCommand>::new(mapping_vmo, PENDING_COMMANDS_CAPACITY)
643                .expect("Failed to create receiver");
644
645        let device = fixture.fs().device().clone();
646        let service: Arc<dyn mapping::reader::BlockService> =
647            Arc::new(DeviceBlockService { device });
648
649        let port = zx::Port::create();
650        let delivery_queue = zx::Vmo::create(mapping::DELIVERY_VMO_SIZE)
651            .expect("Failed to create delivery_queue VMO");
652        let pager = Arc::new(zx::Pager::create(zx::PagerOptions::empty()).unwrap());
653        let vmo_provider = Arc::new(blob_pager_and_verifier::TestVmoProvider::new(
654            pager.clone(),
655            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
656        ));
657        let delivery_receiver = vmo_fifo::Receiver::<mapping::RawDeliveryCommand>::new(
658            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
659            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
660        )
661        .unwrap();
662        let _delivery_processor = blob_pager_and_verifier::DeliveryQueueProcessor::spawn(
663            delivery_receiver,
664            vmo_provider.clone(),
665            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
666        )
667        .unwrap();
668
669        let verifier = block_server::verifier::Verifier::new(delivery_queue);
670        let files = Arc::new(mapping::Files::new(
671            service,
672            verifier,
673            port.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
674        ));
675
676        let id: [u8; 32] = hash.into();
677        let key = 1;
678        let blob_size = session_proxy
679            .open(key, &id)
680            .await
681            .expect("open failed")
682            .expect("open explicitly returned an error");
683
684        assert_eq!(blob_size as usize, uncompressed_data.len());
685
686        let msg = receiver.peek().expect("Failed to peek message");
687        if let Ok(MappingCommand::Mappings { extent_count, metadata_count, .. }) =
688            MappingCommand::try_from(*msg)
689        {
690            if (extent_count as usize) < min_data_extents
691                || (metadata_count as usize) < min_merkle_extents
692            {
693                panic!(
694                    "EXTENTS MISMATCH: extent_count = {}, metadata_count = {}, \
695                     required min_data = {}, min_merkle = {}",
696                    extent_count, metadata_count, min_data_extents, min_merkle_extents
697                );
698            }
699        }
700        mapping::process_mapping_command(&msg, &files).expect("process_mapping_command failed");
701        msg.pop().expect("pop failed");
702
703        let paged_vmo = pager.create_vmo(zx::VmoOptions::empty(), &port, key, blob_size).unwrap();
704        vmo_provider
705            .register_vmo(key, paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap());
706
707        let _pager_thread = files.spawn_pager_thread();
708
709        let (tx, rx) = oneshot::channel();
710        let len = uncompressed_data.len();
711        std::thread::spawn(move || {
712            let mut buf = vec![0u8; len];
713            paged_vmo.read(&mut buf, 0).expect("paged vmo read failed");
714            let _ = tx.send(buf);
715        });
716
717        let read_bytes = rx.await.unwrap();
718        assert_eq!(&read_bytes[..], &uncompressed_data[..]);
719
720        session_proxy.close(key).await.expect("close failed").expect("close error");
721        let msg = receiver.peek().expect("Failed to peek close message");
722        mapping::process_mapping_command(&msg, &files).expect("process_mapping_command failed");
723        msg.pop().expect("pop failed");
724    }
725
726    async fn fragment_free_space(
727        fixture: TestFixture,
728        block_size: usize,
729        anchor_size: Option<usize>,
730    ) -> TestFixture {
731        let mut data = vec![0u8; block_size];
732        let mut hashes = Vec::new();
733        let mut i = 0usize;
734        loop {
735            rand::fill(&mut data[..]);
736            data[..8].copy_from_slice(&(i as u64).to_le_bytes());
737            let hash = fuchsia_merkle::root_from_slice(&data);
738            let delivery_data = Type1Blob::generate(&data, CompressionMode::Never);
739            let writer = match fixture.create_blob(&hash.into(), false).await {
740                Ok(w) => w,
741                Err(_) => break,
742            };
743            if let Ok(mut blob_writer) =
744                BlobWriter::create(writer, delivery_data.len() as u64).await
745            {
746                if blob_writer.write(&delivery_data).await.is_err() {
747                    break;
748                }
749            } else {
750                break;
751            }
752            hashes.push(hash);
753            i += 1;
754        }
755
756        // Unlink every second small blob to create scattered free space holes.
757        let root = fixture.root();
758        for ix in (0..hashes.len()).step_by(2) {
759            let _ = root.unlink(&format!("{}", hashes[ix]), &UnlinkOptions::default()).await;
760        }
761
762        // Remount fixture so unlinked blobs are purged and blocks are freed to allocator.
763        let device = fixture.close().await;
764        let fixture = open_blob_fixture(device).await;
765
766        if let Some(size) = anchor_size {
767            let anchor_data = vec![0u8; size];
768            let _ = fixture.write_blob(&anchor_data, CompressionMode::Never).await;
769        }
770
771        fixture
772    }
773
774    #[fuchsia::test]
775    async fn test_fxfs_blob_mapping_uncompressed() {
776        let fixture = new_blob_fixture().await;
777        let test_data = vec![123u8; 8192];
778        run_blob_mapping_test(&fixture, &test_data, CompressionMode::Never, 1, 0).await;
779        fixture.close().await;
780    }
781
782    #[fuchsia::test]
783    async fn test_fxfs_blob_mapping_compressed() {
784        let fixture = new_blob_fixture().await;
785        // Generate pseudo-random compressible test data.
786        let mut test_data = vec![0u8; 16384];
787        for i in 0..test_data.len() {
788            test_data[i] = ((i / 64) % 256) as u8;
789        }
790        run_blob_mapping_test(&fixture, &test_data, CompressionMode::Always, 1, 0).await;
791        fixture.close().await;
792    }
793
794    #[fuchsia::test]
795    async fn test_fxfs_blob_mapping_fragmented_data() {
796        let fixture = new_blob_fixture().await;
797        let fixture = fragment_free_space(fixture, 32768, Some(1_000_000)).await;
798
799        let mut uncompressed_data = vec![0u8; 262_144];
800        rand::fill(&mut uncompressed_data[..]);
801        let hash1 = write_blob_chunked(&fixture, &uncompressed_data, CompressionMode::Never).await;
802        run_blob_mapping_test_with_hash(&fixture, hash1, &uncompressed_data, 2, 0).await;
803        fixture
804            .root()
805            .unlink(&format!("{}", hash1), &UnlinkOptions::default())
806            .await
807            .unwrap()
808            .unwrap();
809
810        let mut compressed_data = vec![0u8; 262_144];
811        rand::fill(&mut compressed_data[..]);
812        let hash2 = write_blob_chunked(&fixture, &compressed_data, CompressionMode::Always).await;
813        run_blob_mapping_test_with_hash(&fixture, hash2, &compressed_data, 2, 0).await;
814
815        fixture.close().await;
816    }
817}