Skip to main content

fxfs_platform_testing/fuchsia/
volumes_directory.rs

1// Copyright 2022 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::RemoteCrypt;
6use crate::fuchsia::component::map_to_raw_status;
7use crate::fuchsia::directory::FxDirectory;
8use crate::fuchsia::errors::map_to_status;
9use crate::fuchsia::fxblob::BlobDirectory;
10use crate::fuchsia::memory_pressure::{MemoryPressureLevel, MemoryPressureMonitor};
11use crate::fuchsia::profile::new_profile_state;
12use crate::fuchsia::volume::{FxVolume, FxVolumeAndRoot, MemoryPressureConfig, RootDir};
13use anyhow::{Context, Error, anyhow, ensure};
14use fidl::endpoints::{DiscoverableProtocolMarker, ServerEnd};
15use fidl_fuchsia_fs::{AdminMarker, AdminRequest, AdminRequestStream};
16use fidl_fuchsia_fs_startup::{
17    CheckOptions, CreateOptions, MountOptions, VolumeRequest, VolumeRequestStream,
18};
19use fidl_fuchsia_fxfs::{FileBackedVolumeProviderMarker, ProjectIdMarker};
20use fidl_fuchsia_io as fio;
21use fs_inspect::{FsInspectTree, FsInspectVolume};
22use fuchsia_async as fasync;
23use futures::stream::FuturesUnordered;
24use futures::{StreamExt, TryStreamExt};
25use fxfs::errors::FxfsError;
26use fxfs::filesystem::FxFilesystem;
27use fxfs::fsck;
28use fxfs::log::*;
29use fxfs::object_store::transaction::{LockKey, Options, lock_keys};
30use fxfs::object_store::volume::RootVolume;
31use fxfs::object_store::{
32    Directory, NewChildStoreOptions, ObjectDescriptor, ObjectStore, StoreOptions,
33};
34use fxfs_crypto::Crypt;
35use fxfs_trace::{TraceFutureExt, trace_future_args};
36use refaults_vmo::PageRefaultCounter;
37use rustc_hash::FxHashMap as HashMap;
38use std::sync::atomic::{AtomicU64, Ordering};
39use std::sync::{Arc, OnceLock, Weak};
40use vfs::directory::entry_container::MutableDirectory;
41use vfs::directory::helper::DirectlyMutable;
42const MEBIBYTE: u64 = 1024 * 1024;
43
44struct ProfileState {
45    task: fasync::Task<()>,
46    profile_name: String,
47    all_volumes: bool,
48}
49
50/// VolumesDirectory is a special pseudo-directory used to enumerate and operate on volumes.
51/// Volume creation happens via fuchsia.fs.startup.Volumes.Create, rather than open.
52///
53/// Note that VolumesDirectory assumes exclusive access to |root_volume| and if volumes are
54/// manipulated from elsewhere, strange things will happen.
55pub struct VolumesDirectory {
56    root_volume: RootVolume,
57    directory_node: Arc<vfs::directory::immutable::Simple>,
58    mounted_volumes: futures::lock::Mutex<HashMap<u64, MountedVolume>>,
59    inspect_tree: Weak<FsInspectTree>,
60    mem_monitor: Option<MemoryPressureMonitor>,
61    blob_resupplied_count: Arc<PageRefaultCounter>,
62    // The state of profile recordings. Should be locked *after* mounted_volumes.
63    profiling_state: futures::lock::Mutex<Option<ProfileState>>,
64
65    /// A running estimate of the number of dirty bytes outstanding in all pager-backed VMOs across
66    /// all volumes.
67    pager_dirty_bytes_count: PagerDirtyByteCount,
68
69    /// Max outstanding dirty bytes under critical memory pressure. This could be hardcoded, but is
70    /// broken out for testing.
71    max_dirty_bytes_when_critical: AtomicU64,
72
73    // A callback to invoke when a volume is added.  When the volume is removed, this is called
74    // again with `None` as the second parameter.
75    on_volume_added:
76        OnceLock<Box<dyn Fn(&str, Option<(Arc<FxVolume>, Arc<ObjectStore>)>) + Send + Sync>>,
77
78    /// The cache configuration to use under different memory pressure levels.
79    memory_pressure_config: MemoryPressureConfig,
80}
81
82/// Operations on VolumesDirectory that cannot be performed concurrently (i.e. most
83/// volume creation/removal ops) should exist on this guard instead of VolumesDirectory.
84pub struct MountedVolumesGuard<'a> {
85    volumes_directory: Arc<VolumesDirectory>,
86    mounted_volumes: futures::lock::MutexGuard<'a, HashMap<u64, MountedVolume>>,
87}
88
89struct MountedVolume {
90    sequence: u64,
91    volume: FxVolumeAndRoot,
92
93    // True if the volume was forcibly locked.
94    locked: bool,
95}
96
97pub(crate) enum Mode {
98    Mount,
99    Create { guid: Option<[u8; 16]>, low_32_bit_object_ids: bool },
100}
101
102impl MountedVolumesGuard<'_> {
103    /// Creates or mounts a volume. If |crypt| is set, the volume will be created or mounted as
104    /// encrypted. The volume is mounted according to |as_blob|.
105    async fn create_or_mount_volume(
106        &mut self,
107        name: &str,
108        crypt: Option<Arc<dyn Crypt>>,
109        mode: Mode,
110        as_blob: bool,
111    ) -> Result<FxVolumeAndRoot, Error> {
112        let store = match mode {
113            Mode::Create { guid, low_32_bit_object_ids } => self
114                .volumes_directory
115                .root_volume
116                .new_volume(
117                    name,
118                    NewChildStoreOptions {
119                        options: StoreOptions { crypt },
120                        guid,
121                        low_32_bit_object_ids,
122                        ..Default::default()
123                    },
124                )
125                .await
126                .context("failed to create new volume")?,
127            Mode::Mount => {
128                self.volumes_directory.root_volume.volume(name, StoreOptions { crypt }).await?
129            }
130        };
131        ensure!(
132            !self.mounted_volumes.contains_key(&store.store_object_id()),
133            FxfsError::AlreadyBound
134        );
135
136        let volume = if as_blob {
137            self.mount_store::<BlobDirectory>(
138                name,
139                store,
140                self.volumes_directory.memory_pressure_config,
141            )
142            .await?
143        } else {
144            self.mount_store::<FxDirectory>(
145                name,
146                store,
147                self.volumes_directory.memory_pressure_config,
148            )
149            .await?
150        };
151        // If there is an ongoing profile activity, we should apply it to the mounted volume.
152        if let Some(ProfileState { profile_name, all_volumes: true, .. }) =
153            &(*self.volumes_directory.profiling_state.lock().await)
154        {
155            if let Err(e) = volume
156                .volume()
157                .record_and_replay_profile(new_profile_state(as_blob), profile_name)
158                .await
159            {
160                error!(
161                    "Failed to record or replay profile '{}' for volume {}: {:?}",
162                    profile_name, name, e
163                );
164            }
165        }
166
167        if let Mode::Create { .. } = mode {
168            let store_object_id = volume.volume().store().store_object_id();
169            self.volumes_directory.add_directory_entry(name, store_object_id);
170        }
171        Ok(volume)
172    }
173
174    /// Returns the volume if it is found and unlocked, along with a bool to indicate that it is or
175    /// isn't a blob volume.
176    async fn get_unlocked_volume_by_name(
177        &self,
178        volume_name: &str,
179    ) -> Result<(FxVolumeAndRoot, bool), zx::Status> {
180        let (store_object_id, _, _) = self
181            .volumes_directory
182            .root_volume
183            .volume_directory()
184            .lookup(volume_name)
185            .await
186            .map_err(map_to_status)?
187            .ok_or(zx::Status::NOT_FOUND)?;
188        if let Some(MountedVolume { volume, .. }) = self.mounted_volumes.get(&store_object_id) {
189            let is_blob = volume.root().clone().into_any().downcast::<BlobDirectory>().is_ok();
190            Ok((volume.clone(), is_blob))
191        } else {
192            Err(zx::Status::UNAVAILABLE)
193        }
194    }
195
196    // Verifies the integrity of `name` if it exists.  Returns Ok otherwise.
197    async fn check_volume(
198        &mut self,
199        filesystem: &FxFilesystem,
200        name: &str,
201        crypt: Option<Arc<dyn Crypt>>,
202    ) -> Result<Option<Arc<ObjectStore>>, Error> {
203        if let Ok(store) = self
204            .volumes_directory
205            .root_volume
206            .volume(name, StoreOptions { crypt: crypt.clone() })
207            .await
208        {
209            fsck::fsck_volume(filesystem, store.store_object_id(), crypt).await?;
210            Ok(Some(store))
211        } else {
212            Ok(None)
213        }
214    }
215
216    // Mounts the given store.  A lock *must* be held on the volume directory.
217    async fn mount_store<T: From<Directory<FxVolume>> + RootDir>(
218        &mut self,
219        name: &str,
220        store: Arc<ObjectStore>,
221        flush_task_config: MemoryPressureConfig,
222    ) -> Result<FxVolumeAndRoot, Error> {
223        let unique_id = zx::Event::create();
224        let volume = FxVolumeAndRoot::new::<T>(
225            Arc::downgrade(&self.volumes_directory),
226            store,
227            unique_id.koid().unwrap().raw_koid(),
228            name.to_owned(),
229            self.volumes_directory.blob_resupplied_count.clone(),
230            self.volumes_directory.memory_pressure_config,
231        )
232        .await?;
233        volume
234            .volume()
235            .start_background_task(flush_task_config, self.volumes_directory.mem_monitor.as_ref());
236        self.add_mount(name, &volume);
237        Ok(volume)
238    }
239
240    /// Adds a volume (`FxVolumeAndRoot`) into the mount list.
241    pub fn add_mount(&mut self, name: &str, volume: &FxVolumeAndRoot) {
242        static SEQUENCE: AtomicU64 = AtomicU64::new(0);
243        let sequence = SEQUENCE.fetch_add(1, Ordering::Relaxed);
244        self.mounted_volumes.insert(
245            volume.volume().store().store_object_id(),
246            MountedVolume { sequence, volume: volume.clone(), locked: false },
247        );
248        if let Some(inspect) = self.volumes_directory.inspect_tree.upgrade() {
249            inspect.register_volume(
250                name.to_string(),
251                Arc::downgrade(volume.volume()) as Weak<dyn FsInspectVolume + Send + Sync>,
252            )
253        }
254        if let Some(callback) = self.volumes_directory.on_volume_added.get() {
255            callback(
256                name,
257                Some((
258                    volume.volume().clone(),
259                    self.volumes_directory
260                        .root_volume
261                        .volume_directory()
262                        .store()
263                        .filesystem()
264                        .root_store(),
265                )),
266            );
267        }
268    }
269
270    async fn lock_mount(&self, mounted_volume: &mut MountedVolume) {
271        let MountedVolume { volume, locked, .. } = mounted_volume;
272
273        if let Some(callback) = self.volumes_directory.on_volume_added.get() {
274            callback(&volume.volume().name(), None);
275        }
276        if !*locked {
277            if let Some(inspect) = self.volumes_directory.inspect_tree.upgrade() {
278                inspect.unregister_volume(volume.volume().name());
279            }
280            // We must make sure to remove the root entry which holds strong references to the
281            // volume since otherwise `Volume::try_unwrap` might fail.
282            let _ = volume.outgoing_dir().remove_entry("root", true);
283            volume.volume().terminate().await;
284            *locked = true;
285        }
286    }
287
288    async fn remove_volume(&mut self, name: &str) -> Result<(), Error> {
289        let (object_id, transaction) = self
290            .volumes_directory
291            .root_volume
292            .acquire_transaction_for_remove_volume(name, [], false)
293            .await?;
294
295        // Cowardly refuse to delete a mounted volume.
296        ensure!(!self.mounted_volumes.contains_key(&object_id), FxfsError::AlreadyBound);
297        let directory_node = self.volumes_directory.directory_node.clone();
298        self.volumes_directory
299            .root_volume
300            .delete_volume(name, transaction, || {
301                // This shouldn't fail because the entry should exist.
302                directory_node.remove_entry(name, /* must_be_directory: */ false).unwrap();
303            })
304            .await?;
305        Ok(())
306    }
307
308    async fn terminate(&mut self) {
309        let mut volumes = std::mem::take(&mut *self.mounted_volumes);
310        for mounted_volume in volumes.values_mut() {
311            let admin_scope = mounted_volume.volume.admin_scope();
312            admin_scope.shutdown();
313            admin_scope.wait().await;
314
315            self.lock_mount(mounted_volume).await;
316        }
317    }
318
319    // Unmounts the volume identified by `store_id`.  The caller should take locks to avoid races if
320    // necessary.
321    //
322    // NOTE: This will not terminate any connections on the admin scope.
323    pub async fn unmount(&mut self, store_id: u64) -> Result<FxVolumeAndRoot, Error> {
324        let mut mounted_volume =
325            self.mounted_volumes.remove(&store_id).ok_or(FxfsError::NotFound)?;
326        self.lock_mount(&mut mounted_volume).await;
327        Ok(mounted_volume.volume)
328    }
329
330    // Auto-unmount the volume when the last connection to the volume is closed.
331    fn auto_unmount(&self, store_id: u64) {
332        let volumes_directory = self.volumes_directory.clone();
333        let mounted_volume = self.mounted_volumes.get(&store_id).unwrap();
334        let sequence = mounted_volume.sequence;
335        let admin_scope = mounted_volume.volume.admin_scope().clone();
336        let scope = mounted_volume.volume.volume().scope().clone();
337        fasync::Task::spawn(
338            async move {
339                // Check the admin_scope first because once that has finished, there can never be
340                // any more connections to it.
341                admin_scope.wait().await;
342                scope.wait().await;
343
344                // Check that the same volume is still mounted i.e. there wasn't an explicit
345                // unmount.
346                let mut mounted_volumes = volumes_directory.lock().await;
347                match mounted_volumes.mounted_volumes.get(&store_id) {
348                    Some(m) if m.sequence == sequence => {}
349                    _ => return,
350                }
351
352                warn!(store_id; "Last connection to volume closed without unmount, shutting down");
353                let root_store = volumes_directory.root_volume.volume_directory().store();
354                let fs = root_store.filesystem();
355                let _guard = fs
356                    .lock_manager()
357                    .txn_lock(lock_keys![LockKey::object(
358                        root_store.store_object_id(),
359                        volumes_directory.root_volume.volume_directory().object_id(),
360                    )])
361                    .await;
362
363                if let Err(e) = mounted_volumes.unmount(store_id).await {
364                    warn!(e:?, store_id; "Failed to unmount volume");
365                }
366            }
367            .trace(trace_future_args!("Volume::auto_unmount")),
368        )
369        .detach();
370    }
371}
372
373impl VolumesDirectory {
374    /// Fills the VolumesDirectory with all volumes in |root_volume|.  No volume is opened during
375    /// this.
376    pub async fn new(
377        root_volume: RootVolume,
378        inspect_tree: Weak<FsInspectTree>,
379        mem_monitor: Option<MemoryPressureMonitor>,
380        blob_resupplied_count: Arc<PageRefaultCounter>,
381        memory_pressure_config: MemoryPressureConfig,
382    ) -> Result<Arc<Self>, Error> {
383        let layer_set = root_volume.volume_directory().store().tree().layer_set();
384        let mut merger = layer_set.merger();
385        let me = Arc::new(Self {
386            root_volume,
387            directory_node: vfs::directory::immutable::simple(),
388            mounted_volumes: futures::lock::Mutex::new(HashMap::default()),
389            inspect_tree,
390            mem_monitor,
391            blob_resupplied_count,
392            profiling_state: futures::lock::Mutex::new(None),
393            pager_dirty_bytes_count: PagerDirtyByteCount::new(),
394            max_dirty_bytes_when_critical: AtomicU64::new(zx::system_get_physmem() / 100),
395            on_volume_added: OnceLock::new(),
396            memory_pressure_config,
397        });
398        let mut iter = me.root_volume.volume_directory().iter(&mut merger).await?;
399        while let Some((name, store_id, object_descriptor)) = iter.get() {
400            ensure!(*object_descriptor == ObjectDescriptor::Volume, FxfsError::Inconsistent);
401
402            me.add_directory_entry(&name, store_id);
403
404            iter.advance().await?;
405        }
406        Ok(me)
407    }
408
409    /// Clears the directory entry cache (dirent cache) for all mounted volumes.  This drops the
410    /// strong references to nodes held by the cache, allowing them (and their associated VMOs) to
411    /// be freed if they are not otherwise open.  Note that this does not clear other internal
412    /// caches (like LSM tree caches or CachingObjectHandle).
413    pub async fn clear_caches(&self) {
414        let volumes = self.mounted_volumes.lock().await;
415        for mounted_volume in volumes.values() {
416            mounted_volume.volume.volume().dirent_cache().clear();
417        }
418    }
419
420    /// Delete a profile for a given volume. Fails if that volume isn't mounted or if there is
421    /// active profile recording or replay.
422    pub async fn delete_profile(
423        self: &Arc<Self>,
424        volume_name: &str,
425        profile_name: &str,
426    ) -> Result<(), zx::Status> {
427        // Volumes lock is taken first to provide consistent lock ordering with mounting a volume.
428        let volumes = self.mounted_volumes.lock().await;
429        let state = self.profiling_state.lock().await;
430
431        // Only allow deletion when no operations are in flight. This removes confusion around
432        // deleting a profile while one is recording with the same name, as the profile will not be
433        // available for deletion until the recording completes. This would also mean that deleting
434        // during a recording may succeed for deleting an older version but will be confusingly
435        // replaced moments later.
436        if state.is_some() {
437            warn!("Failing profile deletion while profile operations are in flight.");
438            return Err(zx::Status::SHOULD_WAIT);
439        }
440        for MountedVolume { volume, .. } in volumes.values() {
441            if volume.volume().name() == volume_name {
442                let dir = Arc::new(FxDirectory::new(
443                    None,
444                    volume.volume().get_profile_directory().await.map_err(map_to_status)?,
445                ));
446                return dir.unlink(profile_name, false).await;
447            }
448        }
449        warn!(volume_name, profile_name; "Volume not found while deleting profile");
450        Err(zx::Status::NOT_FOUND)
451    }
452
453    pub fn memory_pressure_monitor(&self) -> Option<&MemoryPressureMonitor> {
454        self.mem_monitor.as_ref()
455    }
456
457    /// Stop all ongoing replays, and complete and persist ongoing recordings.
458    pub async fn stop_profile_tasks(self: &Arc<Self>) {
459        let mut state;
460        let volumes;
461        // Take the mounted_volumes lock first to keep consistent lock ordering with other
462        // operations, but don't need to hold it for the entire operation. We need to take the
463        // profiling_state lock before the mounted_volumes lock is dropped to ensure that another
464        // thread doesn't mount a volume and start a profile task on it in between.
465        {
466            volumes = self
467                .mounted_volumes
468                .lock()
469                .await
470                .values()
471                .map(|v| v.volume.volume().clone()) // Clones of each FxVolume.
472                .collect::<Vec<Arc<FxVolume>>>();
473            state = self.profiling_state.lock().await;
474        }
475        for volume in volumes {
476            volume.stop_profile_tasks().await;
477        }
478        *state = None;
479    }
480
481    /// Record a named profile for a number of seconds, fails if there is an in flight recording or
482    /// replay. The given volume must be unlocked, if no volume is given then all volumes will be
483    /// affected and all volumes mounted during the process will also be affected.
484    pub async fn record_and_replay_profile(
485        self: &Arc<Self>,
486        volume_name: Option<String>,
487        profile_name: String,
488        duration_secs: u32,
489    ) -> Result<(), zx::Status> {
490        // Volumes lock is taken first to provide consistent lock ordering with mounting a volume.
491        let volumes = self.lock().await;
492        let mut state = self.profiling_state.lock().await;
493        if state.is_some() {
494            // Consistency in the recording and replaying cannot be ensured at the volume level
495            // if more than one operation can be in flight at a time.
496            return Err(zx::Status::SHOULD_WAIT);
497        }
498        match volume_name.as_ref() {
499            Some(volume_name) => {
500                let (volume, is_blob) = volumes.get_unlocked_volume_by_name(&volume_name).await?;
501                if let Err(error) = volume
502                    .volume()
503                    .record_and_replay_profile(new_profile_state(is_blob), &profile_name)
504                    .await
505                {
506                    error!(
507                        error:?,
508                        profile_name = profile_name.as_str(),
509                        volume_name = volume_name.as_str();
510                        "Failed to record or replay profile",
511                    );
512                    return Err(map_to_status(error));
513                }
514            }
515            None => {
516                for MountedVolume { volume, .. } in volumes.mounted_volumes.values() {
517                    let is_blob =
518                        volume.root().clone().into_any().downcast::<BlobDirectory>().is_ok();
519                    // Just log the errors, don't stop half-way.
520                    if let Err(error) = volume
521                        .volume()
522                        .record_and_replay_profile(new_profile_state(is_blob), &profile_name)
523                        .await
524                    {
525                        error!(
526                            error:?,
527                            profile_name = profile_name.as_str(),
528                            volume_name = volume.volume().name();
529                            "Failed to record or replay profile",
530                        );
531                    }
532                }
533            }
534        }
535
536        let this = self.clone();
537        let task = fasync::Task::spawn(async move {
538            fasync::Timer::new(fasync::MonotonicDuration::from_seconds(duration_secs.into())).await;
539            this.stop_profile_tasks().await;
540        });
541        *state = Some(ProfileState { task, profile_name, all_volumes: volume_name.is_none() });
542        Ok(())
543    }
544
545    /// Replays a profile if one exists, and only records if one does not exist.
546    /// The given volume must be unlocked.
547    pub async fn replay_xor_record_profile(
548        self: &Arc<Self>,
549        volume_name: String,
550        profile_name: String,
551        duration_secs: u32,
552    ) -> Result<(), zx::Status> {
553        // Volumes lock is taken first to provide consistent lock ordering with mounting a volume.
554        let volumes = self.lock().await;
555        let mut state = self.profiling_state.lock().await;
556        if state.is_some() {
557            return Err(zx::Status::SHOULD_WAIT);
558        }
559        let (volume, is_blob) = volumes.get_unlocked_volume_by_name(&volume_name).await?;
560        if let Err(error) = volume
561            .volume()
562            .replay_xor_record_profile(new_profile_state(is_blob), &profile_name)
563            .await
564        {
565            error!(
566                error:?,
567                profile_name = profile_name.as_str(),
568                volume_name = volume_name.as_str();
569                "Failed to replay or record profile",
570            );
571            return Err(map_to_status(error));
572        }
573
574        let this = self.clone();
575        let task = fasync::Task::spawn(async move {
576            fasync::Timer::new(fasync::MonotonicDuration::from_seconds(duration_secs.into())).await;
577            this.stop_profile_tasks().await;
578        });
579        *state = Some(ProfileState { task, profile_name, all_volumes: false });
580        Ok(())
581    }
582
583    /// Returns the directory node which can be used to provide connections for e.g. enumerating
584    /// entries in the VolumesDirectory.
585    /// Directly manipulating the entries in this node will result in strange behaviour.
586    pub fn directory_node(&self) -> &Arc<vfs::directory::immutable::Simple> {
587        &self.directory_node
588    }
589
590    // This serves as an exclusive lock for operations that manipulate the set of mounted volumes.
591    pub async fn lock<'a>(self: &'a Arc<Self>) -> MountedVolumesGuard<'a> {
592        MountedVolumesGuard {
593            volumes_directory: self.clone(),
594            mounted_volumes: self.mounted_volumes.lock().await,
595        }
596    }
597
598    fn add_directory_entry(self: &Arc<Self>, name: &str, store_id: u64) {
599        let weak = Arc::downgrade(self);
600        let name_owned = Arc::new(name.to_string());
601        self.directory_node
602            .add_entry(
603                name,
604                vfs::service::host(move |requests| {
605                    let weak = weak.clone();
606                    let name = name_owned.clone();
607                    async move {
608                        if let Some(me) = weak.upgrade() {
609                            let _ =
610                                me.handle_volume_requests(name.as_ref(), requests, store_id).await;
611                        }
612                    }
613                }),
614            )
615            .unwrap();
616    }
617
618    /// Creates and mounts a new volume. If |crypt| is set, the volume will be encrypted. The
619    /// volume is mounted according to |as_blob|.
620    pub async fn create_and_mount_volume(
621        self: &Arc<Self>,
622        name: &str,
623        crypt: Option<Arc<dyn Crypt>>,
624        as_blob: bool,
625        guid: Option<[u8; 16]>,
626    ) -> Result<FxVolumeAndRoot, Error> {
627        self.lock()
628            .await
629            .create_or_mount_volume(
630                name,
631                crypt,
632                Mode::Create { guid, low_32_bit_object_ids: false },
633                as_blob,
634            )
635            .await
636    }
637
638    /// Mounts an existing volume. `crypt` will be used to unlock the volume if provided.
639    /// If `as_blob` is `true`, the volume will be mounted as a blob filesystem, otherwise
640    /// it will be treated as a regular fxfs volume.
641    pub async fn mount_volume(
642        self: &Arc<Self>,
643        name: &str,
644        crypt: Option<Arc<dyn Crypt>>,
645        as_blob: bool,
646    ) -> Result<FxVolumeAndRoot, Error> {
647        self.lock().await.create_or_mount_volume(name, crypt, Mode::Mount, as_blob).await
648    }
649
650    /// Checks the integrity of volume `name`, or succeeds silently if the volume doesn't exist.
651    pub async fn check_volume(
652        self: &Arc<Self>,
653        name: &str,
654        crypt: Option<Arc<dyn Crypt>>,
655    ) -> Result<(), Error> {
656        let fs = self.root_volume.volume_directory().store().filesystem();
657        let store = self.lock().await.check_volume(fs.as_ref(), name, crypt).await?;
658        if let Some(store) = store {
659            if store.is_unlocked() {
660                let _ = store.lock().await;
661            }
662        }
663        Ok(())
664    }
665
666    /// Removes a volume. The volume must exist but encrypted volume keys are not required.
667    pub async fn remove_volume(self: &Arc<Self>, name: &str) -> Result<(), Error> {
668        self.lock().await.remove_volume(name).await
669    }
670
671    /// Terminates all opened volumes.  This will not cancel any profiling that might be taking
672    /// place.
673    pub async fn terminate(self: &Arc<Self>) {
674        // Abort the profiling timer task.
675        let profiling_state = self.profiling_state.lock().await.take();
676        if let Some(state) = profiling_state {
677            state.task.abort().await;
678        }
679        self.lock().await.terminate().await;
680        // TODO(https://fxbug.dev/452935329): Turn this into a real assert.
681        debug_assert!(
682            self.pager_dirty_bytes_count.load() == 0,
683            "Leaked {} dirty bytes.",
684            self.pager_dirty_bytes_count.load()
685        );
686    }
687
688    /// Serves the given volume on `outgoing_dir_server_end`.
689    pub fn serve_volume(
690        self: &Arc<Self>,
691        volume: &FxVolumeAndRoot,
692        outgoing_dir_server_end: ServerEnd<fio::DirectoryMarker>,
693        as_blob: bool,
694    ) -> Result<(), Error> {
695        // A note regarding strong references to `FxVolume`: connections to services here are all on
696        // the admin scope.  When we force lock a volume, we want to keep the admin scope running
697        // but terminate the volume.  When a volume is terminated in this way, we want to ensure
698        // that there are no outstanding strong references to the volume.  With that in mind, the
699        // services here should not hold strong references.  They can hold a strong reference to
700        // VolumesDirectory and find the volume via the mount list, or they can hold a weak
701        // reference to the FxVolume.  Before upgrading the weak reference, an active guard must be
702        // acquired on the volume's scope first.  This ensures that no new strong references are
703        // taken once the volume has commenced termination.  The "root" entry just below is an
704        // exception that does hold a strong reference.  We handle that by making sure we remove
705        // that entry before we call `FxVolume::terminate`.
706
707        let outgoing_dir = volume.outgoing_dir();
708        outgoing_dir.add_entry("root", volume.root().clone().as_directory_entry())?;
709        let svc_dir = vfs::directory::immutable::simple();
710        outgoing_dir.add_entry("svc", svc_dir.clone())?;
711
712        let store_id = volume.volume().store().store_object_id();
713        let me = self.clone();
714        svc_dir.add_entry(
715            AdminMarker::PROTOCOL_NAME,
716            vfs::service::host(move |requests| {
717                let me = me.clone();
718                async move {
719                    let _ = me.handle_admin_requests(requests, store_id).await;
720                }
721            }),
722        )?;
723        let vol_scope = volume.volume().scope().clone();
724        let weak_vol = Arc::downgrade(volume.volume());
725        {
726            let vol_scope = vol_scope.clone();
727            let weak_vol = weak_vol.clone();
728            svc_dir.add_entry(
729                ProjectIdMarker::PROTOCOL_NAME,
730                vfs::service::host(move |requests| {
731                    let weak_vol = weak_vol.clone();
732                    let scope = vol_scope.clone();
733                    async move {
734                        let _ =
735                            FxVolume::handle_project_id_requests(weak_vol, scope, requests).await;
736                    }
737                }),
738            )?;
739        }
740        svc_dir.add_entry(
741            FileBackedVolumeProviderMarker::PROTOCOL_NAME,
742            vfs::service::host(move |requests| {
743                let weak_vol = weak_vol.clone();
744                let scope = vol_scope.clone();
745                async move {
746                    let _ = FxVolume::handle_file_backed_volume_provider_requests(
747                        weak_vol, scope, requests,
748                    )
749                    .await;
750                }
751            }),
752        )?;
753        volume.root().clone().register_additional_volume_services(&svc_dir)?;
754
755        let scope = volume.admin_scope().clone();
756        let mut flags = fio::PERM_READABLE | fio::PERM_WRITABLE;
757        if as_blob {
758            flags |= fio::PERM_EXECUTABLE;
759        }
760        vfs::directory::serve_on(Arc::clone(outgoing_dir), flags, scope, outgoing_dir_server_end);
761
762        info!(
763            store_id;
764            "Serving volume, pager port koid={}",
765            fasync::EHandle::local().port().koid().unwrap().raw_koid()
766        );
767        Ok(())
768    }
769
770    /// Creates and serves the volume with the given name.
771    pub async fn create_and_serve_volume(
772        self: &Arc<Self>,
773        name: &str,
774        outgoing_directory_server_end: ServerEnd<fio::DirectoryMarker>,
775        mount_options: MountOptions,
776        create_options: CreateOptions,
777    ) -> Result<(), Error> {
778        let mut guard = self.lock().await;
779        let crypt =
780            mount_options.crypt.map(|crypt| Arc::new(RemoteCrypt::new(crypt)) as Arc<dyn Crypt>);
781        let as_blob = mount_options.as_blob.unwrap_or(false);
782        let guid = create_options.guid;
783        let low_32_bit_object_ids = create_options.restrict_inode_ids_to_32_bit.unwrap_or(false);
784        let volume = guard
785            .create_or_mount_volume(
786                name,
787                crypt,
788                Mode::Create { guid, low_32_bit_object_ids },
789                as_blob,
790            )
791            .await?;
792        self.serve_volume(&volume, outgoing_directory_server_end, as_blob)
793            .context("failed to serve volume")?;
794        guard.auto_unmount(volume.volume().store().store_object_id());
795        Ok(())
796    }
797
798    async fn handle_volume_requests(
799        self: &Arc<Self>,
800        name: &str,
801        mut requests: VolumeRequestStream,
802        store_id: u64,
803    ) -> Result<(), Error> {
804        while let Some(request) = requests.try_next().await? {
805            match request {
806                VolumeRequest::Check { responder, options } => {
807                    async move {
808                        responder.send(self.handle_check(store_id, options).await.map_err(
809                            |error| {
810                                error!(error:?, store_id; "Failed to check volume");
811                                map_to_raw_status(error)
812                            },
813                        ))
814                    }
815                    .trace(trace_future_args!("Volume::Check"))
816                    .await?;
817                }
818                VolumeRequest::Mount { responder, outgoing_directory, options } => {
819                    async move {
820                        responder.send(
821                            self.handle_mount(name, store_id, outgoing_directory, options)
822                                .await
823                                .map_err(|error| {
824                                    error!(error:?, name, store_id; "Failed to mount volume");
825                                    map_to_raw_status(error)
826                                }),
827                        )
828                    }
829                    .trace(trace_future_args!("Volume::Mount"))
830                    .await?;
831                }
832                VolumeRequest::SetLimit { responder, bytes } => {
833                    async move {
834                        responder.send(self.handle_set_limit(store_id, bytes).await.map_err(
835                            |error| {
836                                error!(error:?, store_id; "Failed to set volume limit");
837                                map_to_raw_status(error)
838                            },
839                        ))
840                    }
841                    .trace(trace_future_args!("Volume::SetLimit"))
842                    .await?;
843                }
844                VolumeRequest::GetLimit { responder } => {
845                    fxfs_trace::duration!("Volume::GetLimit");
846                    responder.send(Ok(self.handle_get_limit(store_id)))?
847                }
848                VolumeRequest::GetInfo { responder } => {
849                    async move {
850                        let result = self.handle_get_info(store_id).await.map(|guid| {
851                            fidl_fuchsia_fs_startup::VolumeInfo {
852                                guid: Some(guid),
853                                ..Default::default()
854                            }
855                        });
856                        match result {
857                            Ok(response) => responder.send(Ok(&response)),
858                            Err(error) => {
859                                error!(error:?, store_id; "Failed to get volume info");
860                                responder.send(Err(map_to_raw_status(error)))
861                            }
862                        }
863                    }
864                    .trace(trace_future_args!("Volume::GetInfo"))
865                    .await?;
866                }
867            }
868        }
869        Ok(())
870    }
871
872    pub fn memory_pressure_config(&self) -> &MemoryPressureConfig {
873        &self.memory_pressure_config
874    }
875
876    fn is_flush_required_to_dirty(&self, byte_count: u64) -> bool {
877        let mem_pressure = self
878            .mem_monitor
879            .as_ref()
880            .map(|mem_monitor| mem_monitor.level())
881            .unwrap_or(MemoryPressureLevel::Normal);
882        if !matches!(mem_pressure, MemoryPressureLevel::Critical) {
883            return false;
884        }
885
886        let total_dirty = self.pager_dirty_bytes_count.load();
887        total_dirty + byte_count >= self.max_dirty_bytes_when_critical.load(Ordering::Relaxed)
888    }
889
890    /// Reports that a certain number of bytes will be dirtied in a pager-backed VMO. If the memory
891    /// pressure level is critical and fxfs has lots of dirty pages then a new task will be spawned
892    /// in `volume` to flush the dirty pages before `mark_dirty` is called. If the memory pressure
893    /// level is not critical then `mark_dirty` will be synchronously called.
894    pub fn report_pager_dirty(
895        self: Arc<Self>,
896        byte_count: u64,
897        volume: Arc<FxVolume>,
898        mark_dirty: impl FnOnce() + Send + 'static,
899    ) {
900        if !self.is_flush_required_to_dirty(byte_count) {
901            self.pager_dirty_bytes_count.fetch_add(byte_count);
902            mark_dirty();
903        } else {
904            volume.spawn(
905                async move {
906                    let volumes = self.mounted_volumes.lock().await;
907
908                    // Re-check the number of outstanding pager dirty bytes because another thread
909                    // could have raced and flushed the volumes first.
910                    if self.is_flush_required_to_dirty(byte_count) {
911                        debug!(
912                            "Flushing all volumes. Memory pressure is critical & dirty pager bytes \
913                            ({} MiB) >= limit ({} MiB)",
914                            self.pager_dirty_bytes_count.load() / MEBIBYTE,
915                            self.max_dirty_bytes_when_critical.load(Ordering::Relaxed) / MEBIBYTE
916                        );
917
918                        let flushes = FuturesUnordered::new();
919                        for MountedVolume { volume, .. } in volumes.values() {
920                            let vol = volume.volume().clone();
921                            flushes.push(async move {
922                                vol.minimize_memory().await;
923                            });
924                        }
925
926                        flushes.collect::<()>().await;
927                    }
928                    self.pager_dirty_bytes_count.fetch_add(byte_count);
929                    mark_dirty();
930                }
931                .trace(trace_future_args!("flush-before-mark-dirty")),
932            )
933        }
934    }
935
936    /// Reports that a certain number of bytes were cleaned in a pager-backed VMO.
937    pub fn report_pager_clean(&self, byte_count: u64) {
938        let prev_dirty = self.pager_dirty_bytes_count.fetch_sub(byte_count);
939        // TODO(https://fxbug.dev/452935329): Turn this into a real assert.
940        debug_assert!(prev_dirty >= byte_count, "Underflowed dirty bytes.");
941
942        if prev_dirty < byte_count {
943            // An unlikely scenario, but if there was an underflow, reset the pager dirty bytes to
944            // zero.
945            self.pager_dirty_bytes_count.store(0);
946        }
947    }
948
949    async fn handle_check(
950        self: &Arc<Self>,
951        store_id: u64,
952        options: CheckOptions,
953    ) -> Result<(), Error> {
954        let fs = self.root_volume.volume_directory().store().filesystem();
955        let crypt = if let Some(crypt) = options.crypt {
956            Some(Arc::new(RemoteCrypt::new(crypt)) as Arc<dyn Crypt>)
957        } else {
958            None
959        };
960        let result = fsck::fsck_volume(fs.as_ref(), store_id, crypt).await?;
961        // TODO(b/311550633): Stash result in inspect.
962        info!(store_id:%; "{result:?}");
963        Ok(())
964    }
965
966    async fn handle_set_limit(self: &Arc<Self>, store_id: u64, bytes: u64) -> Result<(), Error> {
967        let store = self.root_volume.volume_directory().store();
968        let mut transaction = store.new_transaction(lock_keys![], Options::default()).await?;
969        store.filesystem().allocator().set_bytes_limit(&mut transaction, store_id, bytes)?;
970        transaction.commit().await?;
971        Ok(())
972    }
973
974    fn handle_get_limit(self: &Arc<Self>, store_id: u64) -> u64 {
975        let fs = self.root_volume.volume_directory().store().filesystem();
976        fs.allocator().get_owner_bytes_limit(store_id).unwrap_or_default()
977    }
978
979    async fn handle_get_info(self: &Arc<Self>, store_id: u64) -> Result<[u8; 16], Error> {
980        let fs = self.root_volume.volume_directory().store().filesystem();
981        let store =
982            fs.object_manager().store(store_id).ok_or_else(|| anyhow!("Store not found"))?;
983        Ok(store.guid())
984    }
985
986    async fn handle_mount(
987        self: &Arc<Self>,
988        name: &str,
989        store_id: u64,
990        outgoing_directory_server_end: ServerEnd<fio::DirectoryMarker>,
991        options: MountOptions,
992    ) -> Result<(), Error> {
993        info!(name:%, store_id:%, options:?; "Received mount request");
994        let crypt = options.crypt.map(|crypt| Arc::new(RemoteCrypt::new(crypt)) as Arc<dyn Crypt>);
995        let as_blob = options.as_blob.unwrap_or(false);
996        let mut guard = self.lock().await;
997        let volume = guard
998            .create_or_mount_volume(name, crypt, Mode::Mount, as_blob)
999            .await
1000            .context("failed to mount volume")?;
1001        self.serve_volume(&volume, outgoing_directory_server_end, as_blob)
1002            .context("failed to serve volume")?;
1003        guard.auto_unmount(volume.volume().store().store_object_id());
1004        Ok(())
1005    }
1006
1007    async fn handle_admin_requests(
1008        self: &Arc<Self>,
1009        mut stream: AdminRequestStream,
1010        store_id: u64,
1011    ) -> Result<(), Error> {
1012        // If the Admin protocol ever supports more methods, this should change to a while.
1013        if let Some(request) = stream.try_next().await.context("Reading request")? {
1014            match request {
1015                AdminRequest::Shutdown { responder } => {
1016                    info!(store_id; "Received shutdown request for volume");
1017
1018                    let root_store = self.root_volume.volume_directory().store();
1019                    let fs = root_store.filesystem();
1020                    let _guard = fs
1021                        .lock_manager()
1022                        .txn_lock(lock_keys![LockKey::object(
1023                            root_store.store_object_id(),
1024                            self.root_volume.volume_directory().object_id(),
1025                        )])
1026                        .await;
1027
1028                    let maybe_volume = self.lock().await.unmount(store_id).await;
1029                    responder
1030                        .send()
1031                        .unwrap_or_else(|e| warn!("Failed to send shutdown response: {}", e));
1032
1033                    if let Ok(volume) = maybe_volume {
1034                        // NOTE: After calling this, this task might be dropped at the next await
1035                        // point.
1036                        volume.admin_scope().shutdown();
1037                    }
1038
1039                    return Ok(());
1040                }
1041            }
1042        }
1043        Ok(())
1044    }
1045
1046    /// Sets a callback which is invoked when a volume is added.  When the volume is removed, this
1047    /// is called again with `None` as the second parameter. The root store is included along with
1048    /// volume so the volume's layers can be accessed. Note that this can only be set once per
1049    /// VolumesDirectory; repeated calls will panic.
1050    pub fn set_on_mount_callback<
1051        F: Fn(&str, Option<(Arc<FxVolume>, Arc<ObjectStore>)>) + Send + Sync + 'static,
1052    >(
1053        &self,
1054        callback: F,
1055    ) {
1056        self.on_volume_added.set(Box::new(callback)).ok().unwrap();
1057    }
1058
1059    pub async fn install_volume(
1060        self: &Arc<Self>,
1061        src: &str,
1062        image_file: &str,
1063        dst: &str,
1064    ) -> Result<(), Error> {
1065        let guard = self.lock().await;
1066        info!("installing {src}/{image_file} -> {dst}");
1067        for MountedVolume { volume, .. } in guard.mounted_volumes.values() {
1068            if volume.volume().name() == src {
1069                return Err(zx::Status::ALREADY_BOUND)
1070                    .with_context(|| format!("volume {src} is already mounted"));
1071            }
1072            if volume.volume().name() == dst {
1073                return Err(zx::Status::ALREADY_BOUND)
1074                    .with_context(|| format!("volume {dst} is already mounted"));
1075            }
1076        }
1077        guard.volumes_directory.root_volume.install_volume(&src, &image_file, &dst).await?;
1078
1079        // The above function ensures that we've deleted `src` and `dst` now exists. Before we
1080        // release `guard`, we need to update the entries in the volumes directory accordingly.
1081        guard
1082            .volumes_directory
1083            .directory_node()
1084            .remove_entry(src, /* must_be_directory: */ false)
1085            .unwrap();
1086        guard
1087            .volumes_directory
1088            .directory_node()
1089            .remove_entry(dst, /* must_be_directory: */ false)
1090            .unwrap();
1091        let new_dst_object_id =
1092            match guard.volumes_directory.root_volume.volume_directory().lookup(dst).await? {
1093                Some((object_id, ObjectDescriptor::Volume, _)) => Ok(object_id),
1094                Some(_) => Err(FxfsError::Inconsistent),
1095                None => Err(FxfsError::NotFound),
1096            }?;
1097        self.add_directory_entry(dst, new_dst_object_id);
1098
1099        info!("install complete");
1100        Ok(())
1101    }
1102}
1103
1104#[cfg(test)]
1105pub(crate) fn serve_startup_volume_proxy(
1106    volumes_directory: &Arc<VolumesDirectory>,
1107    volume_name: &str,
1108) -> (fidl_fuchsia_fs_startup::VolumeProxy, vfs::ExecutionScope) {
1109    use vfs::ToObjectRequest;
1110    use vfs::service::ServiceLike;
1111    let scope = vfs::ExecutionScope::new();
1112    let entry = volumes_directory.directory_node().get_entry(volume_name).unwrap();
1113    let service = entry.into_any().downcast::<vfs::service::Service>().unwrap();
1114    let (proxy, server) = fidl::endpoints::create_proxy::<fidl_fuchsia_fs_startup::VolumeMarker>();
1115    service
1116        .connect(
1117            scope.clone(),
1118            Default::default(),
1119            &mut fio::Flags::PROTOCOL_SERVICE.to_object_request(server),
1120        )
1121        .unwrap();
1122    (proxy, scope)
1123}
1124
1125struct PagerDirtyByteCount(AtomicU64);
1126
1127impl PagerDirtyByteCount {
1128    pub fn new() -> Self {
1129        Self(AtomicU64::new(0))
1130    }
1131
1132    pub fn fetch_add(&self, value: u64) -> u64 {
1133        let prev = self.0.fetch_add(value, Ordering::Relaxed);
1134        fxfs_trace::counter!("dirty-bytes", 0, "total" => prev.saturating_add(value));
1135        prev
1136    }
1137
1138    pub fn fetch_sub(&self, value: u64) -> u64 {
1139        let prev = self.0.fetch_sub(value, Ordering::Relaxed);
1140        fxfs_trace::counter!("dirty-bytes", 0, "total" => prev.saturating_sub(value));
1141        prev
1142    }
1143
1144    pub fn load(&self) -> u64 {
1145        self.0.load(Ordering::Relaxed)
1146    }
1147
1148    pub fn store(&self, value: u64) {
1149        self.0.store(value, Ordering::Relaxed);
1150        fxfs_trace::counter!("dirty-bytes", 0, "total" => value);
1151    }
1152}
1153
1154#[cfg(test)]
1155mod tests {
1156    use super::{Mode, serve_startup_volume_proxy};
1157    use crate::fuchsia::RemoteCrypt;
1158    use crate::fuchsia::memory_pressure::MemoryPressureLevel;
1159    use crate::fuchsia::testing::{self, TestFixture, open_dir_checked, open_file_checked};
1160    use crate::fuchsia::volume::MemoryPressureConfig;
1161    use crate::fuchsia::volumes_directory::VolumesDirectory;
1162    use crate::testing::TestFixtureOptions;
1163    use fidl::endpoints::{DiscoverableProtocolMarker, create_proxy, create_request_stream};
1164    use fidl_fuchsia_fs::AdminMarker;
1165    use fidl_fuchsia_fs_startup::{MountOptions, VolumeProxy};
1166    use fidl_fuchsia_fxfs::{CryptRequest, FxfsKey, KeyPurpose, WrappedKey};
1167    use fidl_fuchsia_io as fio;
1168    use fuchsia_async as fasync;
1169    use fuchsia_component_client::connect_to_protocol_at_dir_svc;
1170    use fuchsia_fs::file;
1171    use futures::{TryStreamExt, join};
1172    use fxfs::errors::FxfsError;
1173    use fxfs::filesystem::FxFilesystem;
1174    use fxfs::fsck::{FsckOptions, fsck, fsck_volume_with_options, fsck_with_options};
1175    use fxfs::lock_keys;
1176    use fxfs::object_handle::ObjectHandle;
1177    use fxfs::object_store::allocator::Allocator;
1178    use fxfs::object_store::transaction::{LockKey, Options};
1179    use fxfs::object_store::volume::root_volume;
1180    use fxfs_crypto::Crypt;
1181    use fxfs_insecure_crypto::new_insecure_crypt;
1182    use refaults_vmo::PageRefaultCounter;
1183    use std::sync::atomic::Ordering;
1184    use std::sync::{Arc, Weak};
1185    use std::time::Duration;
1186    use storage_device::DeviceHolder;
1187    use storage_device::fake_device::FakeDevice;
1188    use vfs::execution_scope::ExecutionScope;
1189    use vfs::temp_clone::{TempClonable, unblock};
1190    use zx::Status;
1191    async fn write_image_to_file(image: DeviceHolder, file: fio::FileProxy) {
1192        file.resize(image.size()).await.unwrap().expect("resize failed");
1193        let vmo = TempClonable::new(
1194            file.get_backing_memory(fio::VmoFlags::SHARED_BUFFER | fio::VmoFlags::WRITE)
1195                .await
1196                .unwrap()
1197                .expect("get backing memory failed"),
1198        );
1199
1200        const CHUNK_READ_SIZE: usize = 131_072; /* 128 KiB */
1201        let mut buff = image.allocate_buffer(CHUNK_READ_SIZE).await;
1202        let total = image.size();
1203        let mut offset = 0;
1204        while offset < total {
1205            let amount = std::cmp::min(total - offset, CHUNK_READ_SIZE as u64);
1206            image.read(offset, buff.as_mut()).await.expect("image read failed");
1207            {
1208                // *NOTE*: We have to unblock our write to the VMO since it's pager backed and could
1209                // be running on the same thread as the filesystem is.
1210                let vmo = vmo.temp_clone();
1211                let data = buff.as_slice()[0..amount as usize].to_vec();
1212                let offset = offset;
1213                unblock(move || vmo.write(&data, offset)).await.expect("vmo write failed");
1214            }
1215            offset += amount;
1216        }
1217        assert_eq!(offset, total);
1218        file.sync().await.unwrap().expect("sync failed");
1219        file.close().await.unwrap().expect("close failed");
1220    }
1221
1222    #[fuchsia::test]
1223    async fn test_volume_creation() {
1224        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1225        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1226        let blob_resupplied_count =
1227            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1228        let volumes_directory = VolumesDirectory::new(
1229            root_volume(filesystem.clone()).await.unwrap(),
1230            Weak::new(),
1231            None,
1232            blob_resupplied_count,
1233            MemoryPressureConfig::default(),
1234        )
1235        .await
1236        .unwrap();
1237
1238        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1239        {
1240            let vol = volumes_directory
1241                .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1242                .await
1243                .expect("create encrypted volume failed");
1244            vol.volume().store().store_object_id()
1245        };
1246
1247        volumes_directory.terminate().await;
1248        std::mem::drop(volumes_directory);
1249        filesystem.close().await.expect("close filesystem failed");
1250        let device = filesystem.take_device().await;
1251        device.reopen(false);
1252        let filesystem = FxFilesystem::open(device).await.unwrap();
1253        let blob_resupplied_count =
1254            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1255        let volumes_directory = VolumesDirectory::new(
1256            root_volume(filesystem.clone()).await.unwrap(),
1257            Weak::new(),
1258            None,
1259            blob_resupplied_count,
1260            MemoryPressureConfig::default(),
1261        )
1262        .await
1263        .unwrap();
1264
1265        let error = volumes_directory
1266            .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1267            .await
1268            .err()
1269            .expect("Creating existing encrypted volume should fail");
1270        assert!(FxfsError::AlreadyExists.matches(&error));
1271    }
1272
1273    #[fuchsia::test]
1274    async fn test_dirty_pages_accumulate_in_parent() {
1275        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1276        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1277        let blob_resupplied_count =
1278            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1279        let volumes_directory = VolumesDirectory::new(
1280            root_volume(filesystem.clone()).await.unwrap(),
1281            Weak::new(),
1282            None,
1283            blob_resupplied_count,
1284            MemoryPressureConfig::default(),
1285        )
1286        .await
1287        .unwrap();
1288
1289        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1290        let vol = volumes_directory
1291            .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1292            .await
1293            .expect("create encrypted volume failed");
1294        let old_dirty = volumes_directory.pager_dirty_bytes_count.load();
1295
1296        let new_dirty = {
1297            let (root, server_end) = create_proxy::<fio::DirectoryMarker>();
1298            vol.root().clone().serve(fio::PERM_READABLE | fio::PERM_WRITABLE, server_end);
1299            let f = open_file_checked(
1300                &root,
1301                "foo",
1302                fio::Flags::FLAG_MAYBE_CREATE
1303                    | fio::PERM_READABLE
1304                    | fio::PERM_WRITABLE
1305                    | fio::Flags::PROTOCOL_FILE,
1306                &Default::default(),
1307            )
1308            .await;
1309            let buf = vec![0xaa as u8; 8192];
1310            file::write(&f, buf.as_slice()).await.expect("Write");
1311            // It's important to check the dirty bytes before closing the file, as closing can
1312            // trigger a flush.
1313            volumes_directory.pager_dirty_bytes_count.load()
1314        };
1315        assert_ne!(old_dirty, new_dirty);
1316
1317        volumes_directory.terminate().await;
1318        std::mem::drop(volumes_directory);
1319        filesystem.close().await.expect("close filesystem failed");
1320    }
1321
1322    #[fuchsia::test]
1323    async fn test_volume_reopen() {
1324        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1325        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1326        let blob_resupplied_count =
1327            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1328        let volumes_directory = VolumesDirectory::new(
1329            root_volume(filesystem.clone()).await.unwrap(),
1330            Weak::new(),
1331            None,
1332            blob_resupplied_count,
1333            MemoryPressureConfig::default(),
1334        )
1335        .await
1336        .unwrap();
1337
1338        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1339        let volume_id = {
1340            let vol = volumes_directory
1341                .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1342                .await
1343                .expect("create encrypted volume failed");
1344            vol.volume().store().store_object_id()
1345        };
1346
1347        volumes_directory.terminate().await;
1348        std::mem::drop(volumes_directory);
1349        filesystem.close().await.expect("close filesystem failed");
1350        let device = filesystem.take_device().await;
1351        device.reopen(false);
1352        let filesystem = FxFilesystem::open(device).await.unwrap();
1353        let blob_resupplied_count =
1354            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1355        let volumes_directory = VolumesDirectory::new(
1356            root_volume(filesystem.clone()).await.unwrap(),
1357            Weak::new(),
1358            None,
1359            blob_resupplied_count,
1360            MemoryPressureConfig::default(),
1361        )
1362        .await
1363        .unwrap();
1364
1365        {
1366            let vol = volumes_directory
1367                .mount_volume("encrypted", Some(crypt.clone()), false)
1368                .await
1369                .expect("open existing encrypted volume failed");
1370            assert_eq!(vol.volume().store().store_object_id(), volume_id);
1371        }
1372
1373        volumes_directory.terminate().await;
1374        std::mem::drop(volumes_directory);
1375        filesystem.close().await.expect("close filesystem failed");
1376    }
1377
1378    #[fuchsia::test]
1379    async fn test_volume_creation_unencrypted() {
1380        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1381        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1382        let blob_resupplied_count =
1383            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1384        let volumes_directory = VolumesDirectory::new(
1385            root_volume(filesystem.clone()).await.unwrap(),
1386            Weak::new(),
1387            None,
1388            blob_resupplied_count,
1389            MemoryPressureConfig::default(),
1390        )
1391        .await
1392        .unwrap();
1393
1394        {
1395            let vol = volumes_directory
1396                .create_and_mount_volume("unencrypted", None, false, None)
1397                .await
1398                .expect("create unencrypted volume failed");
1399            vol.volume().store().store_object_id()
1400        };
1401
1402        volumes_directory.terminate().await;
1403        std::mem::drop(volumes_directory);
1404        filesystem.close().await.expect("close filesystem failed");
1405        let device = filesystem.take_device().await;
1406        device.reopen(false);
1407        let filesystem = FxFilesystem::open(device).await.unwrap();
1408        let blob_resupplied_count =
1409            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1410        let volumes_directory = VolumesDirectory::new(
1411            root_volume(filesystem.clone()).await.unwrap(),
1412            Weak::new(),
1413            None,
1414            blob_resupplied_count,
1415            MemoryPressureConfig::default(),
1416        )
1417        .await
1418        .unwrap();
1419
1420        let error = volumes_directory
1421            .create_and_mount_volume("unencrypted", None, false, None)
1422            .await
1423            .err()
1424            .expect("Creating existing unencrypted volume should fail");
1425        assert!(FxfsError::AlreadyExists.matches(&error));
1426
1427        volumes_directory.terminate().await;
1428        std::mem::drop(volumes_directory);
1429        filesystem.close().await.expect("close filesystem failed");
1430    }
1431
1432    #[fuchsia::test]
1433    async fn test_volume_reopen_unencrypted() {
1434        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1435        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1436        let blob_resupplied_count =
1437            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1438        let volumes_directory = VolumesDirectory::new(
1439            root_volume(filesystem.clone()).await.unwrap(),
1440            Weak::new(),
1441            None,
1442            blob_resupplied_count,
1443            MemoryPressureConfig::default(),
1444        )
1445        .await
1446        .unwrap();
1447
1448        let volume_id = {
1449            let vol = volumes_directory
1450                .create_and_mount_volume("unencrypted", None, false, None)
1451                .await
1452                .expect("create unencrypted volume failed");
1453            vol.volume().store().store_object_id()
1454        };
1455
1456        volumes_directory.terminate().await;
1457        std::mem::drop(volumes_directory);
1458        filesystem.close().await.expect("close filesystem failed");
1459        let device = filesystem.take_device().await;
1460        device.reopen(false);
1461        let filesystem = FxFilesystem::open(device).await.unwrap();
1462        let blob_resupplied_count =
1463            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1464        let volumes_directory = VolumesDirectory::new(
1465            root_volume(filesystem.clone()).await.unwrap(),
1466            Weak::new(),
1467            None,
1468            blob_resupplied_count,
1469            MemoryPressureConfig::default(),
1470        )
1471        .await
1472        .unwrap();
1473
1474        {
1475            let vol = volumes_directory
1476                .mount_volume("unencrypted", None, false)
1477                .await
1478                .expect("open existing unencrypted volume failed");
1479            assert_eq!(vol.volume().store().store_object_id(), volume_id);
1480        }
1481
1482        volumes_directory.terminate().await;
1483        std::mem::drop(volumes_directory);
1484        filesystem.close().await.expect("close filesystem failed");
1485    }
1486
1487    #[fuchsia::test]
1488    async fn test_volume_enumeration() {
1489        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1490        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1491        let blob_resupplied_count =
1492            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1493        let volumes_directory = VolumesDirectory::new(
1494            root_volume(filesystem.clone()).await.unwrap(),
1495            Weak::new(),
1496            None,
1497            blob_resupplied_count,
1498            MemoryPressureConfig::default(),
1499        )
1500        .await
1501        .unwrap();
1502
1503        // Add an encrypted volume...
1504        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1505        {
1506            volumes_directory
1507                .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1508                .await
1509                .expect("create encrypted volume failed");
1510        };
1511        // And an unencrypted volume.
1512        {
1513            volumes_directory
1514                .create_and_mount_volume("unencrypted", None, false, None)
1515                .await
1516                .expect("create unencrypted volume failed");
1517        };
1518
1519        // Restart, so that we can test enumeration of unopened volumes.
1520        volumes_directory.terminate().await;
1521        std::mem::drop(volumes_directory);
1522        filesystem.close().await.expect("close filesystem failed");
1523        let device = filesystem.take_device().await;
1524        device.reopen(false);
1525        let filesystem = FxFilesystem::open(device).await.unwrap();
1526        let blob_resupplied_count =
1527            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1528        let volumes_directory = VolumesDirectory::new(
1529            root_volume(filesystem.clone()).await.unwrap(),
1530            Weak::new(),
1531            None,
1532            blob_resupplied_count,
1533            MemoryPressureConfig::default(),
1534        )
1535        .await
1536        .unwrap();
1537
1538        let readdir = |dir: Arc<fio::DirectoryProxy>| async move {
1539            let status = dir.rewind().await.expect("FIDL call failed");
1540            Status::ok(status).expect("rewind failed");
1541            let (status, buf) = dir.read_dirents(fio::MAX_BUF).await.expect("FIDL call failed");
1542            Status::ok(status).expect("read_dirents failed");
1543            let mut entries = vec![];
1544            for res in fuchsia_fs::directory::parse_dir_entries(&buf) {
1545                entries.push(res.expect("Failed to parse entry").name);
1546            }
1547            entries
1548        };
1549
1550        let dir_proxy = Arc::new(vfs::directory::serve_read_only(
1551            volumes_directory.directory_node().clone(),
1552            ExecutionScope::new(),
1553        ));
1554        let entries = readdir(dir_proxy.clone()).await;
1555        assert_eq!(entries, [".", "encrypted", "unencrypted"]);
1556
1557        let _vol = volumes_directory
1558            .mount_volume("encrypted", Some(crypt.clone()), false)
1559            .await
1560            .expect("Open encrypted volume failed");
1561
1562        // Ensure that the behaviour is the same after we've opened a volume.
1563        let entries = readdir(dir_proxy).await;
1564        assert_eq!(entries, [".", "encrypted", "unencrypted"]);
1565
1566        volumes_directory.terminate().await;
1567        std::mem::drop(volumes_directory);
1568        filesystem.close().await.expect("close filesystem failed");
1569    }
1570
1571    #[fuchsia::test]
1572    async fn test_get_info() {
1573        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1574        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1575        let root_volume = root_volume(filesystem.clone()).await.unwrap();
1576        let blob_resupplied_count =
1577            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1578        let volumes_directory = VolumesDirectory::new(
1579            root_volume,
1580            Weak::new(),
1581            None,
1582            blob_resupplied_count,
1583            MemoryPressureConfig::default(),
1584        )
1585        .await
1586        .unwrap();
1587
1588        let vol = volumes_directory
1589            .create_and_mount_volume("vol", None, false, None)
1590            .await
1591            .expect("create_and_mount_volume failed");
1592        let guid = vol.volume().store().guid();
1593
1594        let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, "vol");
1595
1596        let info: fidl_fuchsia_fs_startup::VolumeInfo = volume_proxy
1597            .get_info()
1598            .await
1599            .expect("get_info failed")
1600            .expect("get_info returned error");
1601        assert_eq!(info.guid, Some(guid));
1602
1603        volumes_directory.terminate().await;
1604    }
1605
1606    #[fuchsia::test]
1607    async fn test_deleted_encrypted_volume_while_mounted() {
1608        const VOLUME_NAME: &str = "encrypted";
1609
1610        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1611        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1612        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1613        let blob_resupplied_count =
1614            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1615        let volumes_directory = VolumesDirectory::new(
1616            root_volume(filesystem.clone()).await.unwrap(),
1617            Weak::new(),
1618            None,
1619            blob_resupplied_count,
1620            MemoryPressureConfig::default(),
1621        )
1622        .await
1623        .unwrap();
1624        volumes_directory
1625            .create_and_mount_volume(VOLUME_NAME, Some(crypt.clone()), false, None)
1626            .await
1627            .expect("create encrypted volume failed");
1628        // We have the volume mounted so delete attempts should fail.
1629        assert!(
1630            FxfsError::AlreadyBound.matches(
1631                &volumes_directory
1632                    .remove_volume(VOLUME_NAME)
1633                    .await
1634                    .err()
1635                    .expect("Deleting volume should fail")
1636            )
1637        );
1638        volumes_directory.terminate().await;
1639        std::mem::drop(volumes_directory);
1640        filesystem.close().await.expect("close filesystem failed");
1641    }
1642
1643    #[fuchsia::test]
1644    async fn test_mount_volume_using_volume_protocol() {
1645        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1646        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1647        let blob_resupplied_count =
1648            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1649        let volumes_directory = VolumesDirectory::new(
1650            root_volume(filesystem.clone()).await.unwrap(),
1651            Weak::new(),
1652            None,
1653            blob_resupplied_count,
1654            MemoryPressureConfig::default(),
1655        )
1656        .await
1657        .unwrap();
1658
1659        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1660        let store_id = {
1661            let vol = volumes_directory
1662                .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1663                .await
1664                .expect("create encrypted volume failed");
1665            vol.volume().store().store_object_id()
1666        };
1667        volumes_directory.lock().await.unmount(store_id).await.expect("unmount failed");
1668
1669        let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, "encrypted");
1670
1671        let (dir_proxy, dir_server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1672
1673        let crypt_service = fxfs_crypt::CryptService::new();
1674        crypt_service
1675            .add_wrapping_key(0, fxfs_insecure_crypto::DATA_KEY.to_vec())
1676            .expect("add_wrapping_key failed");
1677        crypt_service
1678            .add_wrapping_key(1, fxfs_insecure_crypto::METADATA_KEY.to_vec())
1679            .expect("add_wrapping_key failed");
1680        crypt_service.set_active_key(KeyPurpose::Data, 0).expect("set_active_key failed");
1681        crypt_service.set_active_key(KeyPurpose::Metadata, 1).expect("set_active_key failed");
1682        let (client1, stream1) = create_request_stream();
1683        let (client2, stream2) = create_request_stream();
1684
1685        join!(
1686            async {
1687                volume_proxy
1688                    .mount(
1689                        dir_server_end,
1690                        MountOptions { crypt: Some(client1), ..MountOptions::default() },
1691                    )
1692                    .await
1693                    .expect("mount (fidl) failed")
1694                    .expect("mount failed");
1695
1696                open_file_checked(
1697                    &dir_proxy,
1698                    "root/test",
1699                    fio::PERM_READABLE | fio::PERM_WRITABLE | fio::Flags::FLAG_MAYBE_CREATE,
1700                    &Default::default(),
1701                )
1702                .await;
1703
1704                // Attempting to mount again should fail with ALREADY_BOUND.
1705                let (_dir_proxy, dir_server_end) =
1706                    fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1707
1708                assert_eq!(
1709                    Status::from_raw(
1710                        volume_proxy
1711                            .mount(
1712                                dir_server_end,
1713                                MountOptions { crypt: Some(client2), ..MountOptions::default() },
1714                            )
1715                            .await
1716                            .expect("mount (fidl) failed")
1717                            .expect_err("mount succeeded")
1718                    ),
1719                    Status::ALREADY_BOUND
1720                );
1721
1722                std::mem::drop(dir_proxy);
1723
1724                // The volume should get unmounted a short time later.
1725                let mut count = 0;
1726                loop {
1727                    if volumes_directory.mounted_volumes.lock().await.is_empty() {
1728                        break;
1729                    }
1730                    count += 1;
1731                    assert!(count <= 100);
1732                    fasync::Timer::new(Duration::from_millis(100)).await;
1733                }
1734            },
1735            async {
1736                crypt_service
1737                    .handle_request(fxfs_crypt::Services::Crypt(stream1))
1738                    .await
1739                    .expect("handle_request failed");
1740                crypt_service
1741                    .handle_request(fxfs_crypt::Services::Crypt(stream2))
1742                    .await
1743                    .expect("handle_request failed");
1744            }
1745        );
1746        // Make sure the background thread that actually calls terminate() on the volume finishes
1747        // before exiting the test. terminate() should be a no-op since we already verified
1748        // mounted_directories is empty, but the volume's terminate() future in the background task
1749        // may still be outstanding. As both the background task and VolumesDirectory::terminate()
1750        // hold the write lock, we use that to block until the background task has completed.
1751        volumes_directory.terminate().await;
1752    }
1753
1754    #[fuchsia::test]
1755    #[ignore] // TODO(b/293917849) re-enable this test when de-flaked
1756
1757    async fn test_volume_dir_races() {
1758        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1759        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1760        let blob_resupplied_count =
1761            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1762        let volumes_directory = VolumesDirectory::new(
1763            root_volume(filesystem.clone()).await.unwrap(),
1764            Weak::new(),
1765            None,
1766            blob_resupplied_count,
1767            MemoryPressureConfig::default(),
1768        )
1769        .await
1770        .unwrap();
1771
1772        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1773        let store_id = {
1774            let vol = volumes_directory
1775                .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1776                .await
1777                .expect("create encrypted volume failed");
1778            vol.volume().store().store_object_id()
1779        };
1780        volumes_directory.lock().await.unmount(store_id).await.expect("unmount failed");
1781
1782        let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, "encrypted");
1783
1784        let crypt_service = Arc::new(fxfs_crypt::CryptService::new());
1785        crypt_service
1786            .add_wrapping_key(0, fxfs_insecure_crypto::DATA_KEY.to_vec())
1787            .expect("add_wrapping_key failed");
1788        crypt_service
1789            .add_wrapping_key(1, fxfs_insecure_crypto::METADATA_KEY.to_vec())
1790            .expect("add_wrapping_key failed");
1791        crypt_service.set_active_key(KeyPurpose::Data, 0).expect("set_active_key failed");
1792        crypt_service.set_active_key(KeyPurpose::Metadata, 1).expect("set_active_key failed");
1793        let (client1, stream1) = create_request_stream();
1794        let (client2, stream2) = create_request_stream();
1795        let crypt_service_clone = crypt_service.clone();
1796        let crypt_task1 = fasync::Task::spawn(async move {
1797            crypt_service_clone
1798                .handle_request(fxfs_crypt::Services::Crypt(stream1))
1799                .await
1800                .expect("handle_request failed");
1801        });
1802        let crypt_task2 = fasync::Task::spawn(async move {
1803            crypt_service
1804                .handle_request(fxfs_crypt::Services::Crypt(stream2))
1805                .await
1806                .expect("handle_request failed");
1807        });
1808
1809        // Create two tasks each of mount and remove, and one to recreate the volume, so that we get
1810        // to exercise a wide variety of concurrent actions.
1811        // Delay remove and create a bit, since mount is slower due to FIDL.
1812        join!(
1813            async {
1814                let (_dir_proxy, dir_server_end) =
1815                    fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1816                if let Err(status) = volume_proxy
1817                    .mount(
1818                        dir_server_end,
1819                        MountOptions { crypt: Some(client1), ..MountOptions::default() },
1820                    )
1821                    .await
1822                    .expect("mount (fidl) failed")
1823                {
1824                    let status = Status::from_raw(status);
1825                    if status != Status::NOT_FOUND && status != Status::ALREADY_BOUND {
1826                        assert!(false, "Unexpected status {:}", status);
1827                    }
1828                }
1829            },
1830            async {
1831                let (_dir_proxy, dir_server_end) =
1832                    fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1833                if let Err(status) = volume_proxy
1834                    .mount(
1835                        dir_server_end,
1836                        MountOptions { crypt: Some(client2), ..MountOptions::default() },
1837                    )
1838                    .await
1839                    .expect("mount (fidl) failed")
1840                {
1841                    let status = Status::from_raw(status);
1842                    if status != Status::NOT_FOUND && status != Status::ALREADY_BOUND {
1843                        assert!(false, "Unexpected status {:}", status);
1844                    }
1845                }
1846            },
1847            async {
1848                let volumes_directory = volumes_directory.clone();
1849                let wait_time = rand::random_range(0..5);
1850                fasync::Timer::new(Duration::from_millis(wait_time)).await;
1851                if let Err(err) = volumes_directory.remove_volume("encrypted").await {
1852                    assert!(
1853                        FxfsError::NotFound.matches(&err) || FxfsError::AlreadyBound.matches(&err),
1854                        "Unexpected error {:?}",
1855                        err
1856                    );
1857                }
1858            },
1859            async {
1860                let volumes_directory = volumes_directory.clone();
1861                let wait_time = rand::random_range(0..5);
1862                fasync::Timer::new(Duration::from_millis(wait_time)).await;
1863                if let Err(err) = volumes_directory.remove_volume("encrypted").await {
1864                    assert!(
1865                        FxfsError::NotFound.matches(&err) || FxfsError::AlreadyBound.matches(&err),
1866                        "Unexpected error {:?}",
1867                        err
1868                    );
1869                }
1870            },
1871            async {
1872                let volumes_directory = volumes_directory.clone();
1873                let wait_time = rand::random_range(0..5);
1874                fasync::Timer::new(Duration::from_millis(wait_time)).await;
1875                let mut guard = volumes_directory.lock().await;
1876                match guard
1877                    .create_or_mount_volume(
1878                        "encrypted",
1879                        Some(crypt.clone()),
1880                        Mode::Create { guid: None, low_32_bit_object_ids: false },
1881                        false,
1882                    )
1883                    .await
1884                {
1885                    Ok(vol) => {
1886                        let store_id = vol.volume().store().store_object_id();
1887                        std::mem::drop(vol);
1888                        guard.unmount(store_id).await.expect("unmount failed");
1889                    }
1890                    Err(err) => {
1891                        assert!(
1892                            FxfsError::AlreadyExists.matches(&err)
1893                                || FxfsError::AlreadyBound.matches(&err),
1894                            "Unexpected error {:?}",
1895                            err
1896                        );
1897                    }
1898                }
1899            }
1900        );
1901        std::mem::drop(crypt_task1);
1902        std::mem::drop(crypt_task2);
1903        // Make sure the background thread that actually calls terminate() on the volume finishes
1904        // before exiting the test. terminate() should be a no-op since we already verified
1905        // mounted_directories is empty, but the volume's terminate() future in the background task
1906        // may still be outstanding. As both the background task and VolumesDirectory::terminate()
1907        // hold the write lock, we use that to block until the background task has completed.
1908        volumes_directory.terminate().await;
1909    }
1910
1911    #[fuchsia::test]
1912    async fn test_shutdown_volume() {
1913        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1914        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1915        let blob_resupplied_count =
1916            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1917        let volumes_directory = VolumesDirectory::new(
1918            root_volume(filesystem.clone()).await.unwrap(),
1919            Weak::new(),
1920            None,
1921            blob_resupplied_count,
1922            MemoryPressureConfig::default(),
1923        )
1924        .await
1925        .unwrap();
1926
1927        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1928        let vol = volumes_directory
1929            .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1930            .await
1931            .expect("create encrypted volume failed");
1932
1933        let (dir_proxy, dir_server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1934
1935        volumes_directory.serve_volume(&vol, dir_server_end, false).expect("serve_volume failed");
1936
1937        let admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(&dir_proxy)
1938            .expect("Unable to connect to admin service");
1939
1940        admin_proxy.shutdown().await.expect("shutdown failed");
1941
1942        assert!(volumes_directory.mounted_volumes.lock().await.is_empty());
1943    }
1944
1945    #[fuchsia::test]
1946    async fn test_byte_limit_persistence() {
1947        const BYTES_LIMIT_1: u64 = 123456;
1948        const BYTES_LIMIT_2: u64 = 456789;
1949        const VOLUME_NAME: &str = "A";
1950        let mut device = DeviceHolder::new(FakeDevice::new(8192, 512));
1951        {
1952            let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1953            let blob_resupplied_count =
1954                Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1955            let volumes_directory = VolumesDirectory::new(
1956                root_volume(filesystem.clone()).await.unwrap(),
1957                Weak::new(),
1958                None,
1959                blob_resupplied_count,
1960                MemoryPressureConfig::default(),
1961            )
1962            .await
1963            .unwrap();
1964
1965            volumes_directory
1966                .create_and_mount_volume(VOLUME_NAME, None, false, None)
1967                .await
1968                .expect("create unencrypted volume failed");
1969
1970            let (volume_proxy, _scope) =
1971                serve_startup_volume_proxy(&volumes_directory, VOLUME_NAME);
1972
1973            volume_proxy.set_limit(BYTES_LIMIT_1).await.unwrap().expect("To set limits");
1974            {
1975                let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
1976                assert_eq!(limits.len(), 1);
1977                assert_eq!(limits[0].1, BYTES_LIMIT_1);
1978            }
1979
1980            volume_proxy.set_limit(BYTES_LIMIT_2).await.unwrap().expect("To set limits");
1981            {
1982                let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
1983                assert_eq!(limits.len(), 1);
1984                assert_eq!(limits[0].1, BYTES_LIMIT_2);
1985            }
1986            std::mem::drop(volume_proxy);
1987            volumes_directory.terminate().await;
1988            std::mem::drop(volumes_directory);
1989            filesystem.close().await.expect("close filesystem failed");
1990            device = filesystem.take_device().await;
1991        }
1992        device.ensure_unique();
1993        device.reopen(false);
1994        {
1995            let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
1996            fsck(filesystem.clone()).await.expect("Fsck");
1997            let blob_resupplied_count =
1998                Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1999            let volumes_directory = VolumesDirectory::new(
2000                root_volume(filesystem.clone()).await.unwrap(),
2001                Weak::new(),
2002                None,
2003                blob_resupplied_count,
2004                MemoryPressureConfig::default(),
2005            )
2006            .await
2007            .unwrap();
2008            {
2009                let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
2010                assert_eq!(limits.len(), 1);
2011                assert_eq!(limits[0].1, BYTES_LIMIT_2);
2012            }
2013            volumes_directory.remove_volume(VOLUME_NAME).await.expect("Volume deletion failed");
2014            {
2015                let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
2016                assert_eq!(limits.len(), 0);
2017            }
2018            volumes_directory.terminate().await;
2019            std::mem::drop(volumes_directory);
2020            filesystem.close().await.expect("close filesystem failed");
2021            device = filesystem.take_device().await;
2022        }
2023        device.ensure_unique();
2024        device.reopen(false);
2025        let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
2026        fsck(filesystem.clone()).await.expect("Fsck");
2027        let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
2028        assert_eq!(limits.len(), 0);
2029    }
2030
2031    struct VolumeInfo {
2032        _scope: vfs::ExecutionScope,
2033        volume_proxy: VolumeProxy,
2034        file_proxy: fio::FileProxy,
2035    }
2036
2037    impl VolumeInfo {
2038        async fn new(volumes_directory: &Arc<VolumesDirectory>, name: &'static str) -> Self {
2039            let volume = volumes_directory
2040                .create_and_mount_volume(name, None, false, None)
2041                .await
2042                .expect("create unencrypted volume failed");
2043
2044            let (volume_dir_proxy, dir_server_end) =
2045                fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2046            volumes_directory
2047                .serve_volume(&volume, dir_server_end, false)
2048                .expect("serve_volume failed");
2049
2050            let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, name);
2051
2052            let (root_proxy, root_server_end) =
2053                fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2054            volume_dir_proxy
2055                .open(
2056                    "root",
2057                    fio::PERM_READABLE | fio::PERM_WRITABLE | fio::Flags::PROTOCOL_DIRECTORY,
2058                    &Default::default(),
2059                    root_server_end.into_channel(),
2060                )
2061                .expect("Failed to open volume root");
2062
2063            let file_proxy = open_file_checked(
2064                &root_proxy,
2065                "foo",
2066                fio::Flags::FLAG_MAYBE_CREATE
2067                    | fio::PERM_READABLE
2068                    | fio::PERM_WRITABLE
2069                    | fio::Flags::PROTOCOL_FILE,
2070                &Default::default(),
2071            )
2072            .await;
2073            VolumeInfo { _scope, volume_proxy, file_proxy }
2074        }
2075    }
2076
2077    #[fuchsia::test]
2078    async fn test_limit_bytes() {
2079        const BYTES_LIMIT: u64 = 262_144; // 256KiB
2080        const BLOCK_SIZE: usize = 8192; // 8KiB
2081        let device = DeviceHolder::new(FakeDevice::new(BLOCK_SIZE.try_into().unwrap(), 512));
2082        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2083        let blob_resupplied_count =
2084            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2085        let volumes_directory = VolumesDirectory::new(
2086            root_volume(filesystem.clone()).await.unwrap(),
2087            Weak::new(),
2088            None,
2089            blob_resupplied_count,
2090            MemoryPressureConfig::default(),
2091        )
2092        .await
2093        .unwrap();
2094
2095        let vol = VolumeInfo::new(&volumes_directory, "foo").await;
2096        let old_info = {
2097            let (status, info) = vol.file_proxy.query_filesystem().await.expect("Getting fs info");
2098            assert_eq!(status, zx::Status::OK.into_raw());
2099            let info = info.unwrap();
2100            // With no limit set, the total filesystem size should be returned.
2101            assert!(info.total_bytes > BYTES_LIMIT);
2102            info
2103        };
2104
2105        vol.volume_proxy.set_limit(BYTES_LIMIT).await.unwrap().expect("To set limits");
2106        {
2107            let (status, info) = vol.file_proxy.query_filesystem().await.expect("Getting fs info");
2108            assert!(status == zx::Status::OK.into_raw());
2109            let new_info = info.unwrap();
2110            assert_eq!(new_info.total_bytes, BYTES_LIMIT);
2111            // Now since the limit is the volume limit, the space used should be the volume usage,
2112            // which should be strictly less than the filesystem.
2113            assert!(new_info.used_bytes < old_info.used_bytes);
2114        }
2115
2116        let zeros = vec![0u8; BLOCK_SIZE];
2117        // First write should succeed.
2118        assert_eq!(
2119            <u64 as TryInto<usize>>::try_into(
2120                vol.file_proxy
2121                    .write(&zeros)
2122                    .await
2123                    .expect("Failed Write message")
2124                    .expect("Failed write")
2125            )
2126            .unwrap(),
2127            BLOCK_SIZE
2128        );
2129        // Likely to run out of space before writing the full limit due to overheads.
2130        for _ in (BLOCK_SIZE..BYTES_LIMIT as usize).step_by(BLOCK_SIZE) {
2131            match vol.file_proxy.write(&zeros).await.expect("Failed Write message") {
2132                Err(_) => break,
2133                Ok(b) if b < BLOCK_SIZE.try_into().unwrap() => break,
2134                _ => (),
2135            };
2136        }
2137
2138        // Any further writes should fail with out of space.
2139        assert_eq!(
2140            vol.file_proxy
2141                .write(&zeros)
2142                .await
2143                .expect("Failed write message")
2144                .expect_err("Write should have been limited"),
2145            Status::NO_SPACE.into_raw()
2146        );
2147
2148        // Double the limit and try again. We should have write space again.
2149        vol.volume_proxy.set_limit(BYTES_LIMIT * 2).await.unwrap().expect("To set limits");
2150        assert_eq!(
2151            <u64 as TryInto<usize>>::try_into(
2152                vol.file_proxy
2153                    .write(&zeros)
2154                    .await
2155                    .expect("Failed Write message")
2156                    .expect("Failed write")
2157            )
2158            .unwrap(),
2159            BLOCK_SIZE
2160        );
2161
2162        vol.file_proxy.close().await.unwrap().expect("Failed to close file");
2163        volumes_directory.terminate().await;
2164        std::mem::drop(volumes_directory);
2165        filesystem.close().await.expect("close filesystem failed");
2166    }
2167
2168    #[fuchsia::test]
2169    async fn test_limit_bytes_two_hit_device_limit() {
2170        const BYTES_LIMIT: u64 = 3_145_728; // 3MiB
2171        const BLOCK_SIZE: usize = 8192; // 8KiB
2172        const BLOCK_COUNT: u32 = 512;
2173        let device =
2174            DeviceHolder::new(FakeDevice::new(BLOCK_SIZE.try_into().unwrap(), BLOCK_COUNT));
2175        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2176        let blob_resupplied_count =
2177            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2178        let volumes_directory = VolumesDirectory::new(
2179            root_volume(filesystem.clone()).await.unwrap(),
2180            Weak::new(),
2181            None,
2182            blob_resupplied_count,
2183            MemoryPressureConfig::default(),
2184        )
2185        .await
2186        .unwrap();
2187
2188        let a = VolumeInfo::new(&volumes_directory, "foo").await;
2189        let b = VolumeInfo::new(&volumes_directory, "bar").await;
2190        a.volume_proxy.set_limit(BYTES_LIMIT).await.unwrap().expect("To set limits");
2191        b.volume_proxy.set_limit(BYTES_LIMIT).await.unwrap().expect("To set limits");
2192        let mut a_written: u64 = 0;
2193        let mut b_written: u64 = 0;
2194
2195        // Write chunks of BLOCK_SIZE.
2196        let zeros = vec![0u8; BLOCK_SIZE];
2197
2198        // First write should succeed for both.
2199        assert_eq!(
2200            <u64 as TryInto<usize>>::try_into(
2201                a.file_proxy
2202                    .write(&zeros)
2203                    .await
2204                    .expect("Failed Write message")
2205                    .expect("Failed write")
2206            )
2207            .unwrap(),
2208            BLOCK_SIZE
2209        );
2210        a_written += BLOCK_SIZE as u64;
2211        assert_eq!(
2212            <u64 as TryInto<usize>>::try_into(
2213                b.file_proxy
2214                    .write(&zeros)
2215                    .await
2216                    .expect("Failed Write message")
2217                    .expect("Failed write")
2218            )
2219            .unwrap(),
2220            BLOCK_SIZE
2221        );
2222        b_written += BLOCK_SIZE as u64;
2223
2224        // Likely to run out of space before writing the full limit due to overheads.
2225        for _ in (BLOCK_SIZE..BYTES_LIMIT as usize).step_by(BLOCK_SIZE) {
2226            match a.file_proxy.write(&zeros).await.expect("Failed Write message") {
2227                Err(_) => break,
2228                Ok(bytes) => {
2229                    a_written += bytes;
2230                    if bytes < BLOCK_SIZE.try_into().unwrap() {
2231                        break;
2232                    }
2233                }
2234            };
2235        }
2236        // Any further writes should fail with out of space.
2237        assert_eq!(
2238            a.file_proxy
2239                .write(&zeros)
2240                .await
2241                .expect("Failed write message")
2242                .expect_err("Write should have been limited"),
2243            Status::NO_SPACE.into_raw()
2244        );
2245
2246        // Now write to the second volume. Likely to run out of space before writing the full limit
2247        // due to overheads.
2248        for _ in (BLOCK_SIZE..BYTES_LIMIT as usize).step_by(BLOCK_SIZE) {
2249            match b.file_proxy.write(&zeros).await.expect("Failed Write message") {
2250                Err(_) => break,
2251                Ok(bytes) => {
2252                    b_written += bytes;
2253                    if bytes < BLOCK_SIZE.try_into().unwrap() {
2254                        break;
2255                    }
2256                }
2257            };
2258        }
2259        // Any further writes should fail with out of space.
2260        assert_eq!(
2261            b.file_proxy
2262                .write(&zeros)
2263                .await
2264                .expect("Failed write message")
2265                .expect_err("Write should have been limited"),
2266            Status::NO_SPACE.into_raw()
2267        );
2268
2269        // Second volume should have failed very early.
2270        assert!(BLOCK_SIZE as u64 * BLOCK_COUNT as u64 - BYTES_LIMIT >= b_written);
2271        // First volume should have gotten further.
2272        assert!(BLOCK_SIZE as u64 * BLOCK_COUNT as u64 - BYTES_LIMIT <= a_written);
2273
2274        a.file_proxy.close().await.unwrap().expect("Failed to close file");
2275        b.file_proxy.close().await.unwrap().expect("Failed to close file");
2276        volumes_directory.terminate().await;
2277        std::mem::drop(volumes_directory);
2278        filesystem.close().await.expect("close filesystem failed");
2279    }
2280
2281    #[fuchsia::test(threads = 10)]
2282    async fn test_profile_start() {
2283        const PREMOUNT_BLOB: &str = "premount_blob";
2284        const PREMOUNT_NOBLOB: &str = "premount_noblob";
2285        const LIVE_BLOB: &str = "live_blob";
2286        const LIVE_NOBLOB: &str = "live_noblob";
2287
2288        const RECORDING_NAME: &str = "foo";
2289
2290        let device = {
2291            let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2292            let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2293            let blob_resupplied_count =
2294                Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2295            let volumes_directory = VolumesDirectory::new(
2296                root_volume(filesystem.clone()).await.unwrap(),
2297                Weak::new(),
2298                None,
2299                blob_resupplied_count,
2300                MemoryPressureConfig::default(),
2301            )
2302            .await
2303            .unwrap();
2304            volumes_directory
2305                .create_and_mount_volume(PREMOUNT_BLOB, None, true, None)
2306                .await
2307                .unwrap();
2308            volumes_directory
2309                .create_and_mount_volume(PREMOUNT_NOBLOB, None, false, None)
2310                .await
2311                .unwrap();
2312            volumes_directory.create_and_mount_volume(LIVE_BLOB, None, true, None).await.unwrap();
2313            volumes_directory
2314                .create_and_mount_volume(LIVE_NOBLOB, None, false, None)
2315                .await
2316                .unwrap();
2317
2318            volumes_directory.terminate().await;
2319            std::mem::drop(volumes_directory);
2320            filesystem.close().await.expect("Filesystem close");
2321            filesystem.take_device().await
2322        };
2323
2324        device.ensure_unique();
2325        device.reopen(false);
2326        let device = {
2327            let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
2328            let blob_resupplied_count =
2329                Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2330            let volumes_directory = VolumesDirectory::new(
2331                root_volume(filesystem.clone()).await.unwrap(),
2332                Weak::new(),
2333                None,
2334                blob_resupplied_count,
2335                MemoryPressureConfig::default(),
2336            )
2337            .await
2338            .unwrap();
2339
2340            // Premount two volumes.
2341            let _premount_blob = volumes_directory
2342                .mount_volume(PREMOUNT_BLOB, None, true)
2343                .await
2344                .expect("Reopen volume");
2345            let _premount_noblob = volumes_directory
2346                .mount_volume(PREMOUNT_NOBLOB, None, false)
2347                .await
2348                .expect("Reopen volume");
2349
2350            // Start the recording, let it run a really long time, it doesn't need to end for this
2351            // test. If it does wait this long then it should trigger test timeouts.
2352            volumes_directory
2353                .clone()
2354                .record_and_replay_profile(None, RECORDING_NAME.to_owned(), 600)
2355                .await
2356                .expect("Recording");
2357
2358            // Live mount two volumes.
2359            let _live_blob =
2360                volumes_directory.mount_volume(LIVE_BLOB, None, true).await.expect("Reopen volume");
2361            let _live_noblob = volumes_directory
2362                .mount_volume(LIVE_NOBLOB, None, false)
2363                .await
2364                .expect("Reopen volume");
2365
2366            // Wait for the recordings to finish.
2367            volumes_directory.stop_profile_tasks().await;
2368
2369            volumes_directory.terminate().await;
2370            std::mem::drop(volumes_directory);
2371            filesystem.close().await.expect("Filesystem close");
2372            filesystem.take_device().await
2373        };
2374
2375        device.ensure_unique();
2376        device.reopen(false);
2377        let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
2378        {
2379            let blob_resupplied_count =
2380                Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2381            let volumes_directory = VolumesDirectory::new(
2382                root_volume(filesystem.clone()).await.unwrap(),
2383                Weak::new(),
2384                None,
2385                blob_resupplied_count,
2386                MemoryPressureConfig::default(),
2387            )
2388            .await
2389            .unwrap();
2390
2391            let _premount_blob = volumes_directory
2392                .mount_volume(PREMOUNT_BLOB, None, true)
2393                .await
2394                .expect("Reopen volume");
2395            let _premount_noblob = volumes_directory
2396                .mount_volume(PREMOUNT_NOBLOB, None, false)
2397                .await
2398                .expect("Reopen volume");
2399            let _live_blob =
2400                volumes_directory.mount_volume(LIVE_BLOB, None, true).await.expect("Reopen volume");
2401            let _live_noblob = volumes_directory
2402                .mount_volume(LIVE_NOBLOB, None, false)
2403                .await
2404                .expect("Reopen volume");
2405
2406            // Verify which recordings ran based on the saved recordings.
2407            volumes_directory
2408                .delete_profile(PREMOUNT_BLOB, RECORDING_NAME)
2409                .await
2410                .expect("Finding profile to delete.");
2411            volumes_directory
2412                .delete_profile(PREMOUNT_NOBLOB, RECORDING_NAME)
2413                .await
2414                .expect("Finding profile to delete.");
2415            volumes_directory
2416                .delete_profile(LIVE_BLOB, RECORDING_NAME)
2417                .await
2418                .expect("Finding profile to delete.");
2419            volumes_directory
2420                .delete_profile(LIVE_NOBLOB, RECORDING_NAME)
2421                .await
2422                .expect("Finding profile to delete.");
2423
2424            volumes_directory.terminate().await;
2425        }
2426
2427        filesystem.close().await.expect("Filesystem close");
2428    }
2429
2430    #[fuchsia::test(threads = 10)]
2431    async fn test_profile_stop() {
2432        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2433        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2434        let blob_resupplied_count =
2435            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2436        let volumes_directory = VolumesDirectory::new(
2437            root_volume(filesystem.clone()).await.unwrap(),
2438            Weak::new(),
2439            None,
2440            blob_resupplied_count,
2441            MemoryPressureConfig::default(),
2442        )
2443        .await
2444        .unwrap();
2445        let volume =
2446            volumes_directory.create_and_mount_volume("foo", None, true, None).await.unwrap();
2447
2448        // Run the recording with no time at all and ensure that it still shuts down properly.
2449        volumes_directory
2450            .clone()
2451            .record_and_replay_profile(None, "foo".to_owned(), 0)
2452            .await
2453            .expect("Recording");
2454
2455        // Delete will succeed once the profile is completed.
2456        while volumes_directory.delete_profile("foo", "foo").await.is_err() {
2457            fasync::Timer::new(Duration::from_millis(10)).await;
2458        }
2459
2460        std::mem::drop(volume);
2461        volumes_directory.terminate().await;
2462        std::mem::drop(volumes_directory);
2463        filesystem.close().await.expect("Filesystem close");
2464    }
2465
2466    #[fuchsia::test(threads = 10)]
2467    async fn test_delete_profile() {
2468        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2469        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2470        let blob_resupplied_count =
2471            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2472        let volumes_directory = VolumesDirectory::new(
2473            root_volume(filesystem.clone()).await.unwrap(),
2474            Weak::new(),
2475            None,
2476            blob_resupplied_count,
2477            MemoryPressureConfig::default(),
2478        )
2479        .await
2480        .unwrap();
2481        let volume =
2482            volumes_directory.create_and_mount_volume("foo", None, true, None).await.unwrap();
2483
2484        volumes_directory
2485            .clone()
2486            .record_and_replay_profile(None, "foo".to_owned(), 600)
2487            .await
2488            .expect("Recording");
2489
2490        // Deletion fails during in-flight recording.
2491        assert_eq!(
2492            volumes_directory.delete_profile("foo", "foo").await.expect_err("File shouldn't exist"),
2493            Status::SHOULD_WAIT
2494        );
2495
2496        volumes_directory.stop_profile_tasks().await;
2497
2498        // Missing volume name.
2499        assert_eq!(
2500            volumes_directory.delete_profile("bar", "foo").await.expect_err("File shouldn't exist"),
2501            Status::NOT_FOUND
2502        );
2503
2504        // Missing Profile name.
2505        assert_eq!(
2506            volumes_directory.delete_profile("foo", "bar").await.expect_err("File shouldn't exist"),
2507            Status::NOT_FOUND
2508        );
2509
2510        // Deletion should now succeed as the profile will be placed as part of `finish_profiling()`
2511        volumes_directory.delete_profile("foo", "foo").await.expect("Deleting");
2512
2513        // Deletion fails as the file shouldn't exist anymore.
2514        assert_eq!(
2515            volumes_directory.delete_profile("foo", "foo").await.expect_err("File shouldn't exist"),
2516            Status::NOT_FOUND
2517        );
2518
2519        std::mem::drop(volume);
2520        volumes_directory.terminate().await;
2521        std::mem::drop(volumes_directory);
2522        filesystem.close().await.expect("Filesystem close");
2523    }
2524
2525    #[fuchsia::test(threads = 10)]
2526    async fn test_profile_start_single_volume() {
2527        const TEST_VOLUME: &str = "test_1234";
2528        const TEST_RECORDING: &str = "test_5678";
2529        let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
2530
2531        let fixture = TestFixture::new().await;
2532        {
2533            let volumes_directory = fixture.volumes_directory();
2534            // Start recording for volume that doesn't exist.
2535            assert_eq!(
2536                volumes_directory
2537                    .record_and_replay_profile(
2538                        Some(TEST_VOLUME.to_owned()),
2539                        TEST_RECORDING.to_owned(),
2540                        1
2541                    )
2542                    .await
2543                    .expect_err("Volumes doesn't exist yet"),
2544                Status::NOT_FOUND
2545            );
2546
2547            // Create the volume unmounted.
2548            {
2549                let volume = volumes_directory
2550                    .create_and_mount_volume(TEST_VOLUME, Some(crypt.clone()), false, None)
2551                    .await
2552                    .unwrap();
2553                volumes_directory
2554                    .lock()
2555                    .await
2556                    .unmount(volume.volume().store().store_object_id())
2557                    .await
2558                    .expect("unmount failed");
2559            }
2560
2561            // Start recording for volume that is not mounted.
2562            assert_eq!(
2563                volumes_directory
2564                    .record_and_replay_profile(
2565                        Some(TEST_VOLUME.to_owned()),
2566                        TEST_RECORDING.to_owned(),
2567                        1
2568                    )
2569                    .await
2570                    .expect_err("Volumes doesn't exist yet"),
2571                Status::UNAVAILABLE
2572            );
2573
2574            // Remount the volume and try again.
2575            let volume = volumes_directory
2576                .mount_volume(TEST_VOLUME, Some(crypt.clone()), false)
2577                .await
2578                .expect("Remount volume");
2579            volumes_directory
2580                .record_and_replay_profile(
2581                    Some(TEST_VOLUME.to_owned()),
2582                    TEST_RECORDING.to_owned(),
2583                    1,
2584                )
2585                .await
2586                .expect("Starting recording");
2587
2588            // Stop the recording and check that it was created on the new volume.
2589            volumes_directory.stop_profile_tasks().await;
2590            {
2591                let profile_dir = volume.volume().get_profile_directory().await.unwrap();
2592                assert!(profile_dir.lookup(TEST_RECORDING).await.unwrap().is_some());
2593            }
2594
2595            // Should be no recording for the other volume.
2596            {
2597                let profile_dir = fixture.volume().volume().get_profile_directory().await.unwrap();
2598                assert!(profile_dir.lookup(TEST_RECORDING).await.unwrap().is_none());
2599            }
2600        }
2601        fixture.close().await;
2602    }
2603
2604    #[fuchsia::test(threads = 10)]
2605    async fn test_delete_volume_while_flushing() {
2606        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2607        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2608        let blob_resupplied_count =
2609            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2610        let volumes_directory = VolumesDirectory::new(
2611            root_volume(filesystem.clone()).await.unwrap(),
2612            Weak::new(),
2613            None,
2614            blob_resupplied_count,
2615            MemoryPressureConfig::default(),
2616        )
2617        .await
2618        .unwrap();
2619        let name = "vol";
2620        let volume =
2621            volumes_directory.create_and_mount_volume(name, None, false, None).await.unwrap();
2622        let mut transaction = filesystem
2623            .root_store()
2624            .new_transaction(
2625                lock_keys![LockKey::object(
2626                    volume.volume().store().store_object_id(),
2627                    volume.root_dir().directory().object_id()
2628                )],
2629                Options::default(),
2630            )
2631            .await
2632            .unwrap();
2633        volume
2634            .root_dir()
2635            .directory()
2636            .create_child_file(&mut transaction, "foo")
2637            .await
2638            .expect("create_child_file failed");
2639        transaction.commit().await.expect("commit failed");
2640        volumes_directory
2641            .lock()
2642            .await
2643            .unmount(volume.volume().store().store_object_id())
2644            .await
2645            .expect("unmount failed");
2646
2647        let filesystem_clone = filesystem.clone();
2648        let filesystem_clone2 = filesystem.clone();
2649        let volumes_directory_clone1 = volumes_directory.clone();
2650        let volumes_directory_clone2 = volumes_directory.clone();
2651        let root_store_object_id = filesystem.root_store().store_object_id();
2652        let store_info_object_id = volume.volume().store().store_info_handle_object_id().unwrap();
2653        join!(
2654            async move {
2655                // Take a lock that the EndFlush transaction requires, so we interleave removing
2656                // the volume between StartFlush and EndFlush.  Release it on a timer since once the
2657                // volume is deleted, `remove_volume` needs to take this lock as well to tombstone
2658                // the store info.
2659                let _guard = filesystem_clone2
2660                    .lock_manager()
2661                    .read_lock(lock_keys![LockKey::object(
2662                        root_store_object_id,
2663                        store_info_object_id,
2664                    )])
2665                    .await;
2666                fasync::Timer::new(Duration::from_millis(200)).await;
2667            },
2668            async move {
2669                filesystem_clone.journal().force_compact().await.expect("Compact failed");
2670            },
2671            async move {
2672                if let Err(e) = volumes_directory_clone1.remove_volume(name).await {
2673                    if !FxfsError::NotFound.matches(&e) {
2674                        panic!("remove_volume failed: {e:?}");
2675                    }
2676                }
2677            },
2678            async move {
2679                if let Err(e) = volumes_directory_clone2.remove_volume(name).await {
2680                    if !FxfsError::NotFound.matches(&e) {
2681                        panic!("remove_volume failed: {e:?}");
2682                    }
2683                }
2684            },
2685        );
2686        volumes_directory.terminate().await;
2687        std::mem::drop(volumes_directory);
2688        filesystem.close().await.expect("Filesystem close");
2689    }
2690
2691    // This mostly just ensures that we exercise the code path. It was added because the path at
2692    // one point contained a deadlock.
2693    #[fuchsia::test(threads = 10)]
2694    async fn test_flush_before_mark_dirty_under_critical_memory_pressure() {
2695        let fixture = TestFixture::new().await;
2696        // Memory is critical, and it's always our fault.
2697        let _ = fixture
2698            .memory_pressure_proxy()
2699            .on_level_changed(MemoryPressureLevel::Critical)
2700            .await
2701            .expect("memory pressure FIDL");
2702        fixture.volumes_directory().max_dirty_bytes_when_critical.store(1, Ordering::Relaxed);
2703
2704        let root = fixture.root();
2705        let file = open_file_checked(
2706            &root,
2707            "foo",
2708            fio::Flags::FLAG_MAYBE_CREATE
2709                | fio::PERM_READABLE
2710                | fio::PERM_WRITABLE
2711                | fio::Flags::PROTOCOL_FILE,
2712            &Default::default(),
2713        )
2714        .await;
2715
2716        file.resize((zx::system_get_page_size() * 2).into())
2717            .await
2718            .expect("resize (FIDL)")
2719            .expect("resize failed");
2720        // The above resize creates zero pages which aren't dirty pages and don't contribute to the
2721        // dirty bytes count but still need to be flushed. The first write below will cross the
2722        // `max_dirty_bytes_when_critical` threshold but `minimize_memory` won't flush the file
2723        // because the zero pages aren't dirty pages. The second write below will also cross the
2724        // `max_dirty_bytes_when_critical` threshold and since the first write dirtied a page, a
2725        // flush will occur. The flush cleans the 2nd page which is a zero page and the kernel is
2726        // actively trying to dirty. The kernel doesn't end up dirtying the 2nd page because it's
2727        // concerned that the clean and dirty raced so it issues another dirty request for the same
2728        // page. The duplicate dirty request again crosses the `max_dirty_bytes_when_critical`
2729        // threshold. The file thinks that it has a dirty page so `minimize_memory` flushes the
2730        // file. `zx_pager_query_dirty_ranges` doesn't return any pages because the kernel didn't
2731        // actually dirty the page that fxfs thought it did. This causes only the file metadata to
2732        // be flushed and the dirty page count to not be reduced. The 2nd page finally gets dirtied
2733        // and now fxfs thinks it has 2 dirty pages instead of 1. `PagedObjectHandle` knows that
2734        // this can happen and will fix up the counts on the next flush. This test doesn't want to
2735        // deal with all of that so it syncs here to clean the zero pages that would cause that
2736        // mess.
2737        file.sync().await.expect("Failed to make sync call").expect("sync failed");
2738
2739        let vmo = file
2740            .get_backing_memory(fio::VmoFlags::READ | fio::VmoFlags::WRITE)
2741            .await
2742            .expect("get_backing_memory (FIDL)")
2743            .expect("get_backing_memory");
2744
2745        let buf = [0xAAu8];
2746        // One call to get dirty bytes over 0, the second to force a flush during mark_dirty.
2747        vmo.write(&buf, 0).expect("Writing to create dirty bytes");
2748        let before = fixture.volumes_directory().pager_dirty_bytes_count.load();
2749        vmo.write(&buf, zx::system_get_page_size().into())
2750            .expect("Writing to force a flush during mark_dirty");
2751        // This is still the page size because we forced a flush of the first write during the
2752        // second write.
2753        assert_eq!(fixture.volumes_directory().pager_dirty_bytes_count.load(), before,);
2754
2755        fixture.close().await;
2756    }
2757
2758    #[fuchsia::test(threads = 10)]
2759    async fn test_delete_crypt_for_volume() {
2760        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2761        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2762        let store_id;
2763        {
2764            let blob_resupplied_count =
2765                Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2766            let volumes_directory = VolumesDirectory::new(
2767                root_volume(filesystem.clone()).await.unwrap(),
2768                Weak::new(),
2769                None,
2770                blob_resupplied_count,
2771                MemoryPressureConfig::default(),
2772            )
2773            .await
2774            .unwrap();
2775            let name = "vol";
2776            let crypt = Arc::new(new_insecure_crypt());
2777            let volume = volumes_directory
2778                .create_and_mount_volume(name, Some(crypt.clone()), false, None)
2779                .await
2780                .unwrap();
2781            store_id = volume.volume().store().store_object_id();
2782            // Make sure the volume has some journaled mutations.
2783            let mut transaction = filesystem
2784                .root_store()
2785                .new_transaction(
2786                    lock_keys![LockKey::object(
2787                        volume.volume().store().store_object_id(),
2788                        volume.root_dir().directory().object_id()
2789                    )],
2790                    Options::default(),
2791                )
2792                .await
2793                .unwrap();
2794            volume
2795                .root_dir()
2796                .directory()
2797                .create_child_file(&mut transaction, "foo")
2798                .await
2799                .expect("create_child_file failed");
2800            transaction.commit().await.expect("commit failed");
2801
2802            let (volume_dir_proxy, dir_server_end) =
2803                fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2804            volumes_directory
2805                .serve_volume(&volume, dir_server_end, false)
2806                .expect("serve_volume failed");
2807            let (root_dir, root_server_end) =
2808                fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2809            volume_dir_proxy
2810                .open(
2811                    "root",
2812                    fio::PERM_READABLE | fio::PERM_WRITABLE | fio::Flags::PROTOCOL_DIRECTORY,
2813                    &Default::default(),
2814                    root_server_end.into_channel(),
2815                )
2816                .expect("Failed to open volume root");
2817
2818            let filesystem_clone = filesystem.clone();
2819            join!(
2820                async move {
2821                    filesystem_clone.journal().force_compact().await.expect("Compact failed");
2822                },
2823                async move {
2824                    let mut i = 0;
2825                    while let Ok(_) = fuchsia_fs::directory::open_file(
2826                        &root_dir,
2827                        &format!("foo{i}"),
2828                        fio::Flags::FLAG_MAYBE_CREATE | fio::PERM_READABLE,
2829                    )
2830                    .await
2831                    {
2832                        i += 1;
2833                    }
2834                },
2835                async move {
2836                    crypt.shutdown();
2837                },
2838            );
2839
2840            // Make sure we can still ask the volume to shutdown.
2841            let (admin_proxy, server_end) =
2842                fidl::endpoints::create_proxy::<fidl_fuchsia_fs::AdminMarker>();
2843            volume_dir_proxy
2844                .open(
2845                    &format!("svc/{}", fidl_fuchsia_fs::AdminMarker::PROTOCOL_NAME),
2846                    fio::Flags::PROTOCOL_SERVICE,
2847                    &Default::default(),
2848                    server_end.into(),
2849                )
2850                .expect("Failed to open Admin connection");
2851            admin_proxy.shutdown().await.expect("shutdown failed");
2852
2853            volumes_directory.terminate().await;
2854        }
2855        filesystem.close().await.expect("Filesystem close");
2856        let device = filesystem.take_device().await;
2857        device.reopen(false);
2858        let filesystem = FxFilesystem::open(device).await.expect("open failed");
2859        let options = FsckOptions { fail_on_warning: true, ..Default::default() };
2860        fsck_with_options(filesystem.clone(), &options).await.expect("fsck failed");
2861        fsck_volume_with_options(
2862            filesystem.as_ref(),
2863            &options,
2864            store_id,
2865            Some(Arc::new(new_insecure_crypt())),
2866        )
2867        .await
2868        .expect("fsck_volume failed");
2869        filesystem.close().await.expect("Filesystem close");
2870    }
2871
2872    /// Tests writing a partition image and installing a volume contained within.
2873    #[fuchsia::test(threads = 10)]
2874    async fn test_volume_installation() {
2875        let fixture = TestFixture::open(
2876            DeviceHolder::new(FakeDevice::new(1024, 4096)),
2877            TestFixtureOptions { format: true, encrypted: false, ..Default::default() },
2878        )
2879        .await;
2880
2881        // Create a file "foo" in the existing volume "vol". This should be gone after installation.
2882        {
2883            let file = open_file_checked(
2884                fixture.root(),
2885                "foo",
2886                fio::Flags::PROTOCOL_FILE | fio::Flags::FLAG_MUST_CREATE | fio::PERM_WRITABLE,
2887                &Default::default(),
2888            )
2889            .await;
2890            file.write("Hello, world!".as_bytes()).await.unwrap().expect("write failed");
2891        };
2892
2893        // Create another in-memory partition image with a different set of files.
2894        let image = {
2895            let inner_fixture = TestFixture::open(
2896                DeviceHolder::new(FakeDevice::new(512, 4096)),
2897                TestFixtureOptions { format: true, encrypted: false, ..Default::default() },
2898            )
2899            .await;
2900            let file = open_file_checked(
2901                inner_fixture.root(),
2902                "bar",
2903                fio::Flags::PROTOCOL_FILE | fio::Flags::FLAG_MUST_CREATE | fio::PERM_WRITABLE,
2904                &Default::default(),
2905            )
2906            .await;
2907            file.write("Well, this is new...".as_bytes()).await.unwrap().expect("write failed");
2908            file.close().await.unwrap().expect("close error");
2909            inner_fixture.close().await
2910        };
2911
2912        // Write the partition image to a file "install" in a new volume called "src".
2913        {
2914            let (src_out_dir, server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2915            fixture
2916                .volumes_directory()
2917                .create_and_serve_volume("src", server_end, Default::default(), Default::default())
2918                .await
2919                .unwrap();
2920            let src_root = open_dir_checked(
2921                &src_out_dir,
2922                "root",
2923                fio::PERM_READABLE | fio::PERM_WRITABLE,
2924                Default::default(),
2925            )
2926            .await;
2927            let file = open_file_checked(
2928                &src_root,
2929                "image",
2930                fio::Flags::PROTOCOL_FILE | fio::Flags::FLAG_MUST_CREATE | fio::PERM_WRITABLE,
2931                &Default::default(),
2932            )
2933            .await;
2934            write_image_to_file(image, file).await;
2935        };
2936
2937        // Installation should not be possible yet since both volumes are still mounted.
2938        assert!(
2939            fixture.volumes_directory().install_volume("src", "image", "vol").await.is_err(),
2940            "volume installation should fail while either src/dst is still mounted"
2941        );
2942
2943        // Now let's re-mount the filesystem manually without a fixture, install the volumes, and
2944        // spin up a new fixture to verify the result.
2945        let device = fixture.close().await;
2946        let fs = FxFilesystem::open(device).await.unwrap();
2947        {
2948            let root = root_volume(fs.clone()).await.unwrap();
2949            root.install_volume("src", "image", "vol").await.unwrap();
2950        }
2951        fs.close().await.unwrap();
2952        let device = fs.take_device().await;
2953        device.reopen(/*read_only*/ true);
2954        let fixture = TestFixture::open(
2955            device,
2956            TestFixtureOptions { encrypted: false, format: false, ..Default::default() },
2957        )
2958        .await;
2959
2960        // Ensure that the "src" volume is now gone, and the old file "foo" is gone from "vol".
2961        assert!(
2962            fixture.volumes_directory().mount_volume("src", None, false).await.is_err(),
2963            "src volume should be deleted after installation"
2964        );
2965        assert!(
2966            testing::open_file(
2967                fixture.volume_out_dir(),
2968                "foo",
2969                fio::PERM_READABLE,
2970                &Default::default()
2971            )
2972            .await
2973            .is_err(),
2974            "foo should be deleted after installation"
2975        );
2976
2977        // Check that we can find the contents of "vol" that we installed from the image.
2978        let file =
2979            open_file_checked(fixture.root(), "bar", fio::PERM_READABLE, &Default::default()).await;
2980        let data = file.read(fio::MAX_TRANSFER_SIZE).await.unwrap().expect("read failed");
2981        assert_eq!(String::from_utf8(data).unwrap(), "Well, this is new...");
2982        file.close().await.unwrap().unwrap();
2983
2984        fixture.close().await;
2985    }
2986
2987    #[fuchsia::test]
2988    async fn test_create_with_low_32_bit_ids() {
2989        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2990        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2991        let blob_resupplied_count =
2992            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2993
2994        {
2995            let volumes_directory = VolumesDirectory::new(
2996                root_volume(filesystem.clone()).await.unwrap(),
2997                Weak::new(),
2998                None,
2999                blob_resupplied_count,
3000                MemoryPressureConfig::default(),
3001            )
3002            .await
3003            .unwrap();
3004
3005            let mut guard = volumes_directory.lock().await;
3006
3007            let vol = guard
3008                .create_or_mount_volume(
3009                    "low_32",
3010                    None,
3011                    Mode::Create { guid: None, low_32_bit_object_ids: true },
3012                    false,
3013                )
3014                .await
3015                .expect("create volume failed");
3016
3017            let root_dir = vol.volume().store().root_directory_object_id();
3018            let root_dir = fxfs::object_store::Directory::open(vol.volume().store(), root_dir)
3019                .await
3020                .expect("open failed");
3021
3022            let mut transaction = filesystem
3023                .root_store()
3024                .new_transaction(
3025                    lock_keys![LockKey::object(
3026                        vol.volume().store().store_object_id(),
3027                        root_dir.object_id()
3028                    )],
3029                    Options::default(),
3030                )
3031                .await
3032                .expect("new_transaction failed");
3033
3034            let object = root_dir
3035                .create_child_file(&mut transaction, "test")
3036                .await
3037                .expect("create_child_file failed");
3038
3039            // We can't check LastObjectIdInfo as it is private, but we can verify behavior.
3040            assert!(object.object_id() < 1 << 32);
3041            transaction.commit().await.expect("commit failed");
3042        };
3043
3044        filesystem.close().await.expect("close filesystem failed");
3045
3046        // Reopen and verify persistence
3047        let device = filesystem.take_device().await;
3048        device.reopen(false);
3049        let filesystem = FxFilesystem::open(device).await.unwrap();
3050        let blob_resupplied_count =
3051            Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
3052        let volumes_directory = VolumesDirectory::new(
3053            root_volume(filesystem.clone()).await.unwrap(),
3054            Weak::new(),
3055            None,
3056            blob_resupplied_count,
3057            MemoryPressureConfig::default(),
3058        )
3059        .await
3060        .unwrap();
3061
3062        // Mount the volume again, and check that new files are still created with expected IDs.
3063        {
3064            let mut guard = volumes_directory.lock().await;
3065
3066            let vol = guard
3067                .create_or_mount_volume("low_32", None, Mode::Mount, false)
3068                .await
3069                .expect("mount volume failed");
3070
3071            let root_dir = vol.volume().store().root_directory_object_id();
3072            let root_dir = fxfs::object_store::Directory::open(vol.volume().store(), root_dir)
3073                .await
3074                .expect("open failed");
3075
3076            let mut transaction = filesystem
3077                .root_store()
3078                .new_transaction(
3079                    lock_keys![LockKey::object(
3080                        vol.volume().store().store_object_id(),
3081                        root_dir.object_id()
3082                    )],
3083                    Options::default(),
3084                )
3085                .await
3086                .expect("new_transaction failed");
3087
3088            let object = root_dir
3089                .create_child_file(&mut transaction, "test2")
3090                .await
3091                .expect("create_child_file failed");
3092
3093            assert!(object.object_id() < 1 << 32);
3094            transaction.commit().await.expect("commit failed");
3095        }
3096
3097        filesystem.close().await.expect("close filesystem failed");
3098    }
3099
3100    #[fuchsia::test(threads = 10)]
3101    async fn test_race_unmount_and_flush_with_crypt_error() {
3102        let device = DeviceHolder::new(FakeDevice::new(8192, 512));
3103        let filesystem = FxFilesystem::new_empty(device).await.unwrap();
3104        let blob_resupplied_count = Arc::new(PageRefaultCounter::new().unwrap());
3105        let volumes_directory = VolumesDirectory::new(
3106            root_volume(filesystem.clone()).await.unwrap(),
3107            Weak::new(),
3108            None,
3109            blob_resupplied_count,
3110            MemoryPressureConfig::default(),
3111        )
3112        .await
3113        .unwrap();
3114
3115        let crypt_service = Arc::new(fxfs_crypt::CryptService::new());
3116        crypt_service
3117            .add_wrapping_key(0, fxfs_insecure_crypto::DATA_KEY.to_vec())
3118            .expect("add_wrapping_key failed");
3119        crypt_service
3120            .add_wrapping_key(1, fxfs_insecure_crypto::METADATA_KEY.to_vec())
3121            .expect("add_wrapping_key failed");
3122        crypt_service.set_active_key(KeyPurpose::Data, 0).expect("set_active_key failed");
3123        crypt_service.set_active_key(KeyPurpose::Metadata, 1).expect("set_active_key failed");
3124
3125        for _ in 0..20 {
3126            let (client, mut stream) = create_request_stream::<fidl_fuchsia_fxfs::CryptMarker>();
3127            let (close_tx, mut close_rx) = futures::channel::oneshot::channel::<()>();
3128
3129            let crypt_task = fasync::Task::spawn(async move {
3130                loop {
3131                    futures::select! {
3132                        _ = close_rx => return, // Close signal
3133                        request = stream.try_next() => {
3134                            match request {
3135                                Ok(Some(CryptRequest::CreateKey { responder, .. })) => {
3136                                    responder.send(Ok((&[0; 16], &[0; 48], &[0; 32]))).unwrap();
3137                                }
3138                                Ok(Some(CryptRequest::CreateKeyWithId {
3139                                        wrapping_key_id, responder, .. })) => {
3140                                    let key = WrappedKey::Fxfs(FxfsKey {
3141                                        wrapping_key_id,
3142                                        wrapped_key: [0u8; 48],
3143                                    });
3144                                    responder.send(Ok((&key, &[0; 32]))).unwrap();
3145                                }
3146                                Ok(Some(CryptRequest::UnwrapKey { responder, .. })) => {
3147                                    responder.send(Ok(&vec![0; 32])).unwrap();
3148                                }
3149                                _ => return,
3150                            }
3151                        }
3152                    }
3153                }
3154            });
3155
3156            let volume = volumes_directory
3157                .create_and_mount_volume(
3158                    "encrypted",
3159                    Some(Arc::new(RemoteCrypt::new(client))),
3160                    false,
3161                    None,
3162                )
3163                .await
3164                .unwrap();
3165
3166            // Write some data to dirty the journal.
3167            {
3168                let mut transaction = filesystem
3169                    .root_store()
3170                    .new_transaction(
3171                        lock_keys![LockKey::object(
3172                            volume.volume().store().store_object_id(),
3173                            volume.root_dir().directory().object_id()
3174                        )],
3175                        Options::default(),
3176                    )
3177                    .await
3178                    .unwrap();
3179                volume
3180                    .root_dir()
3181                    .directory()
3182                    .create_child_file(&mut transaction, "foo")
3183                    .await
3184                    .expect("create_child_file failed");
3185                transaction.commit().await.expect("commit failed");
3186            }
3187
3188            let (dir_proxy, dir_server_end) =
3189                fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
3190            volumes_directory.serve_volume(&volume, dir_server_end, false).unwrap();
3191            volumes_directory.lock().await.auto_unmount(volume.volume().store().store_object_id());
3192
3193            // Kill crypt.  The next usage should be in flush.
3194            let _ = close_tx.send(());
3195
3196            let filesystem_clone = filesystem.clone();
3197            let compact_task = fasync::Task::spawn(async move {
3198                // Flush should not fail, because that would close the journal for the rest of the
3199                // filesystem.  Instead, the volume should be force-locked and flushed (which
3200                // doesn't depend on crypt).
3201                filesystem_clone.object_manager().flush().await.expect("flush failed");
3202            });
3203
3204            // Trigger unmount.
3205            std::mem::drop(dir_proxy);
3206
3207            // Wait for volume to be unmounted, so we can clean up for the next iteration.
3208            let store_id = volume.volume().store().store_object_id();
3209            loop {
3210                {
3211                    let guard = volumes_directory.lock().await;
3212                    if !guard.mounted_volumes.contains_key(&store_id) {
3213                        break;
3214                    }
3215                }
3216                fasync::Timer::new(Duration::from_millis(10)).await;
3217            }
3218
3219            join!(compact_task, crypt_task);
3220            volumes_directory.remove_volume("encrypted").await.expect("remove_volume failed");
3221        }
3222        volumes_directory.terminate().await;
3223        filesystem.close().await.expect("close filesystem failed");
3224    }
3225}