Skip to main content

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