1use 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
50pub 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 profiling_state: futures::lock::Mutex<Option<ProfileState>>,
64
65 pager_dirty_bytes_count: PagerDirtyByteCount,
68
69 max_dirty_bytes_when_critical: AtomicU64,
72
73 on_volume_added:
76 OnceLock<Box<dyn Fn(&str, Option<(Arc<FxVolume>, Arc<ObjectStore>)>) + Send + Sync>>,
77
78 memory_pressure_config: MemoryPressureConfig,
80}
81
82pub 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 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 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 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 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 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 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 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 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 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 directory_node.remove_entry(name, 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 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 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 admin_scope.wait().await;
342 scope.wait().await;
343
344 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 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 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 pub async fn delete_profile(
426 self: &Arc<Self>,
427 volume_name: &str,
428 profile_name: &str,
429 ) -> Result<(), zx::Status> {
430 let volumes = self.mounted_volumes.lock().await;
432 let state = self.profiling_state.lock().await;
433
434 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 pub async fn stop_profile_tasks(self: &Arc<Self>) {
462 let mut state;
463 let volumes;
464 {
469 volumes = self
470 .mounted_volumes
471 .lock()
472 .await
473 .values()
474 .map(|v| v.volume.volume().clone()) .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 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 let volumes = self.lock().await;
495 let mut state = self.profiling_state.lock().await;
496 if state.is_some() {
497 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 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 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 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 pub fn directory_node(&self) -> &Arc<vfs::directory::immutable::Simple> {
590 &self.directory_node
591 }
592
593 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 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 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 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 pub async fn remove_volume(self: &Arc<Self>, name: &str) -> Result<(), Error> {
674 self.lock().await.remove_volume(name).await
675 }
676
677 pub async fn terminate(self: &Arc<Self>) {
680 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 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 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 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 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 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 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 pub fn report_pager_clean(&self, byte_count: u64) {
958 let prev_dirty = self.pager_dirty_bytes_count.fetch_sub(byte_count);
959 debug_assert!(prev_dirty >= byte_count, "Underflowed dirty bytes.");
961
962 if prev_dirty < byte_count {
963 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 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 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 volume.admin_scope().shutdown();
1057 }
1058
1059 return Ok(());
1060 }
1061 }
1062 }
1063 Ok(())
1064 }
1065
1066 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 guard
1102 .volumes_directory
1103 .directory_node()
1104 .remove_entry(src, false)
1105 .unwrap();
1106 guard
1107 .volumes_directory
1108 .directory_node()
1109 .remove_entry(dst, 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; 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 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 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 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 {
1559 volumes_directory
1560 .create_and_mount_volume("unencrypted", None, false, CreateOptions::default())
1561 .await
1562 .expect("create unencrypted volume failed");
1563 };
1564
1565 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 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 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 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 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 volumes_directory.terminate().await;
1808 }
1809
1810 #[fuchsia::test]
1811 #[ignore] 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 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 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 store.flush().await.expect("flush failed");
2066
2067 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 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 assert_eq!(vol.volume().dirent_cache().len(), 0);
2098 assert_eq!(store.tree().cache_len(), 0);
2099 assert!(store.key_manager().get(file_id).await.unwrap().is_none());
2101
2102 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; 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 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 assert!(new_info.used_bytes < old_info.used_bytes);
2284 }
2285
2286 let zeros = vec![0u8; BLOCK_SIZE.get() as usize];
2287 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 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 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 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; 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 let zeros = vec![0u8; BLOCK_SIZE.get() as usize];
2360
2361 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 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 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 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 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 assert!(BLOCK_SIZE * BLOCK_COUNT - BYTES_LIMIT >= b_written);
2420 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 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 volumes_directory
2505 .clone()
2506 .record_and_replay_profile(None, RECORDING_NAME.to_owned(), 600)
2507 .await
2508 .expect("Recording");
2509
2510 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 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 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 volumes_directory
2604 .clone()
2605 .record_and_replay_profile(None, "foo".to_owned(), 0)
2606 .await
2607 .expect("Recording");
2608
2609 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 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 assert_eq!(
2656 volumes_directory.delete_profile("bar", "foo").await.expect_err("File shouldn't exist"),
2657 Status::NOT_FOUND
2658 );
2659
2660 assert_eq!(
2662 volumes_directory.delete_profile("foo", "bar").await.expect_err("File shouldn't exist"),
2663 Status::NOT_FOUND
2664 );
2665
2666 volumes_directory.delete_profile("foo", "foo").await.expect("Deleting");
2668
2669 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 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 {
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 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 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 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 {
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 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 #[fuchsia::test(threads = 10)]
2857 async fn test_flush_before_mark_dirty_under_critical_memory_pressure() {
2858 let fixture = TestFixture::new().await;
2859 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 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 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 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 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 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 #[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 {
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 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 {
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 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 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(true);
3113 let fixture = TestFixture::open(
3114 device,
3115 TestFixtureOptions { encrypted: false, format: false, ..Default::default() },
3116 )
3117 .await;
3118
3119 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 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 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 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 {
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, 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 {
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 let _ = close_tx.send(());
3354
3355 let filesystem_clone = filesystem.clone();
3356 let compact_task = fasync::Task::spawn(async move {
3357 filesystem_clone
3361 .object_manager()
3362 .flush(FlushReason::Journal(ForceMajor::False))
3363 .await
3364 .expect("flush failed");
3365 });
3366
3367 std::mem::drop(dir_proxy);
3369
3370 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}