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::{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) {
414 let volumes = self.mounted_volumes.lock().await;
415 for mounted_volume in volumes.values() {
416 mounted_volume.volume.volume().dirent_cache().clear();
417 }
418 }
419
420 pub async fn delete_profile(
423 self: &Arc<Self>,
424 volume_name: &str,
425 profile_name: &str,
426 ) -> Result<(), zx::Status> {
427 let volumes = self.mounted_volumes.lock().await;
429 let state = self.profiling_state.lock().await;
430
431 if state.is_some() {
437 warn!("Failing profile deletion while profile operations are in flight.");
438 return Err(zx::Status::SHOULD_WAIT);
439 }
440 for MountedVolume { volume, .. } in volumes.values() {
441 if volume.volume().name() == volume_name {
442 let dir = Arc::new(FxDirectory::new(
443 None,
444 volume.volume().get_profile_directory().await.map_err(map_to_status)?,
445 ));
446 return dir.unlink(profile_name, false).await;
447 }
448 }
449 warn!(volume_name, profile_name; "Volume not found while deleting profile");
450 Err(zx::Status::NOT_FOUND)
451 }
452
453 pub fn memory_pressure_monitor(&self) -> Option<&MemoryPressureMonitor> {
454 self.mem_monitor.as_ref()
455 }
456
457 pub async fn stop_profile_tasks(self: &Arc<Self>) {
459 let mut state;
460 let volumes;
461 {
466 volumes = self
467 .mounted_volumes
468 .lock()
469 .await
470 .values()
471 .map(|v| v.volume.volume().clone()) .collect::<Vec<Arc<FxVolume>>>();
473 state = self.profiling_state.lock().await;
474 }
475 for volume in volumes {
476 volume.stop_profile_tasks().await;
477 }
478 *state = None;
479 }
480
481 pub async fn record_and_replay_profile(
485 self: &Arc<Self>,
486 volume_name: Option<String>,
487 profile_name: String,
488 duration_secs: u32,
489 ) -> Result<(), zx::Status> {
490 let volumes = self.lock().await;
492 let mut state = self.profiling_state.lock().await;
493 if state.is_some() {
494 return Err(zx::Status::SHOULD_WAIT);
497 }
498 match volume_name.as_ref() {
499 Some(volume_name) => {
500 let (volume, is_blob) = volumes.get_unlocked_volume_by_name(&volume_name).await?;
501 if let Err(error) = volume
502 .volume()
503 .record_and_replay_profile(new_profile_state(is_blob), &profile_name)
504 .await
505 {
506 error!(
507 error:?,
508 profile_name = profile_name.as_str(),
509 volume_name = volume_name.as_str();
510 "Failed to record or replay profile",
511 );
512 return Err(map_to_status(error));
513 }
514 }
515 None => {
516 for MountedVolume { volume, .. } in volumes.mounted_volumes.values() {
517 let is_blob =
518 volume.root().clone().into_any().downcast::<BlobDirectory>().is_ok();
519 if let Err(error) = volume
521 .volume()
522 .record_and_replay_profile(new_profile_state(is_blob), &profile_name)
523 .await
524 {
525 error!(
526 error:?,
527 profile_name = profile_name.as_str(),
528 volume_name = volume.volume().name();
529 "Failed to record or replay profile",
530 );
531 }
532 }
533 }
534 }
535
536 let this = self.clone();
537 let task = fasync::Task::spawn(async move {
538 fasync::Timer::new(fasync::MonotonicDuration::from_seconds(duration_secs.into())).await;
539 this.stop_profile_tasks().await;
540 });
541 *state = Some(ProfileState { task, profile_name, all_volumes: volume_name.is_none() });
542 Ok(())
543 }
544
545 pub async fn replay_xor_record_profile(
548 self: &Arc<Self>,
549 volume_name: String,
550 profile_name: String,
551 duration_secs: u32,
552 ) -> Result<(), zx::Status> {
553 let volumes = self.lock().await;
555 let mut state = self.profiling_state.lock().await;
556 if state.is_some() {
557 return Err(zx::Status::SHOULD_WAIT);
558 }
559 let (volume, is_blob) = volumes.get_unlocked_volume_by_name(&volume_name).await?;
560 if let Err(error) = volume
561 .volume()
562 .replay_xor_record_profile(new_profile_state(is_blob), &profile_name)
563 .await
564 {
565 error!(
566 error:?,
567 profile_name = profile_name.as_str(),
568 volume_name = volume_name.as_str();
569 "Failed to replay or record profile",
570 );
571 return Err(map_to_status(error));
572 }
573
574 let this = self.clone();
575 let task = fasync::Task::spawn(async move {
576 fasync::Timer::new(fasync::MonotonicDuration::from_seconds(duration_secs.into())).await;
577 this.stop_profile_tasks().await;
578 });
579 *state = Some(ProfileState { task, profile_name, all_volumes: false });
580 Ok(())
581 }
582
583 pub fn directory_node(&self) -> &Arc<vfs::directory::immutable::Simple> {
587 &self.directory_node
588 }
589
590 pub async fn lock<'a>(self: &'a Arc<Self>) -> MountedVolumesGuard<'a> {
592 MountedVolumesGuard {
593 volumes_directory: self.clone(),
594 mounted_volumes: self.mounted_volumes.lock().await,
595 }
596 }
597
598 fn add_directory_entry(self: &Arc<Self>, name: &str, store_id: u64) {
599 let weak = Arc::downgrade(self);
600 let name_owned = Arc::new(name.to_string());
601 self.directory_node
602 .add_entry(
603 name,
604 vfs::service::host(move |requests| {
605 let weak = weak.clone();
606 let name = name_owned.clone();
607 async move {
608 if let Some(me) = weak.upgrade() {
609 let _ =
610 me.handle_volume_requests(name.as_ref(), requests, store_id).await;
611 }
612 }
613 }),
614 )
615 .unwrap();
616 }
617
618 pub async fn create_and_mount_volume(
621 self: &Arc<Self>,
622 name: &str,
623 crypt: Option<Arc<dyn Crypt>>,
624 as_blob: bool,
625 guid: Option<[u8; 16]>,
626 ) -> Result<FxVolumeAndRoot, Error> {
627 self.lock()
628 .await
629 .create_or_mount_volume(
630 name,
631 crypt,
632 Mode::Create { guid, low_32_bit_object_ids: false },
633 as_blob,
634 )
635 .await
636 }
637
638 pub async fn mount_volume(
642 self: &Arc<Self>,
643 name: &str,
644 crypt: Option<Arc<dyn Crypt>>,
645 as_blob: bool,
646 ) -> Result<FxVolumeAndRoot, Error> {
647 self.lock().await.create_or_mount_volume(name, crypt, Mode::Mount, as_blob).await
648 }
649
650 pub async fn check_volume(
652 self: &Arc<Self>,
653 name: &str,
654 crypt: Option<Arc<dyn Crypt>>,
655 ) -> Result<(), Error> {
656 let fs = self.root_volume.volume_directory().store().filesystem();
657 let store = self.lock().await.check_volume(fs.as_ref(), name, crypt).await?;
658 if let Some(store) = store {
659 if store.is_unlocked() {
660 let _ = store.lock().await;
661 }
662 }
663 Ok(())
664 }
665
666 pub async fn remove_volume(self: &Arc<Self>, name: &str) -> Result<(), Error> {
668 self.lock().await.remove_volume(name).await
669 }
670
671 pub async fn terminate(self: &Arc<Self>) {
674 let profiling_state = self.profiling_state.lock().await.take();
676 if let Some(state) = profiling_state {
677 state.task.abort().await;
678 }
679 self.lock().await.terminate().await;
680 debug_assert!(
682 self.pager_dirty_bytes_count.load() == 0,
683 "Leaked {} dirty bytes.",
684 self.pager_dirty_bytes_count.load()
685 );
686 }
687
688 pub fn serve_volume(
690 self: &Arc<Self>,
691 volume: &FxVolumeAndRoot,
692 outgoing_dir_server_end: ServerEnd<fio::DirectoryMarker>,
693 as_blob: bool,
694 ) -> Result<(), Error> {
695 let outgoing_dir = volume.outgoing_dir();
708 outgoing_dir.add_entry("root", volume.root().clone().as_directory_entry())?;
709 let svc_dir = vfs::directory::immutable::simple();
710 outgoing_dir.add_entry("svc", svc_dir.clone())?;
711
712 let store_id = volume.volume().store().store_object_id();
713 let me = self.clone();
714 svc_dir.add_entry(
715 AdminMarker::PROTOCOL_NAME,
716 vfs::service::host(move |requests| {
717 let me = me.clone();
718 async move {
719 let _ = me.handle_admin_requests(requests, store_id).await;
720 }
721 }),
722 )?;
723 let vol_scope = volume.volume().scope().clone();
724 let weak_vol = Arc::downgrade(volume.volume());
725 {
726 let vol_scope = vol_scope.clone();
727 let weak_vol = weak_vol.clone();
728 svc_dir.add_entry(
729 ProjectIdMarker::PROTOCOL_NAME,
730 vfs::service::host(move |requests| {
731 let weak_vol = weak_vol.clone();
732 let scope = vol_scope.clone();
733 async move {
734 let _ =
735 FxVolume::handle_project_id_requests(weak_vol, scope, requests).await;
736 }
737 }),
738 )?;
739 }
740 svc_dir.add_entry(
741 FileBackedVolumeProviderMarker::PROTOCOL_NAME,
742 vfs::service::host(move |requests| {
743 let weak_vol = weak_vol.clone();
744 let scope = vol_scope.clone();
745 async move {
746 let _ = FxVolume::handle_file_backed_volume_provider_requests(
747 weak_vol, scope, requests,
748 )
749 .await;
750 }
751 }),
752 )?;
753 volume.root().clone().register_additional_volume_services(&svc_dir)?;
754
755 let scope = volume.admin_scope().clone();
756 let mut flags = fio::PERM_READABLE | fio::PERM_WRITABLE;
757 if as_blob {
758 flags |= fio::PERM_EXECUTABLE;
759 }
760 vfs::directory::serve_on(Arc::clone(outgoing_dir), flags, scope, outgoing_dir_server_end);
761
762 info!(
763 store_id;
764 "Serving volume, pager port koid={}",
765 fasync::EHandle::local().port().koid().unwrap().raw_koid()
766 );
767 Ok(())
768 }
769
770 pub async fn create_and_serve_volume(
772 self: &Arc<Self>,
773 name: &str,
774 outgoing_directory_server_end: ServerEnd<fio::DirectoryMarker>,
775 mount_options: MountOptions,
776 create_options: CreateOptions,
777 ) -> Result<(), Error> {
778 let mut guard = self.lock().await;
779 let crypt =
780 mount_options.crypt.map(|crypt| Arc::new(RemoteCrypt::new(crypt)) as Arc<dyn Crypt>);
781 let as_blob = mount_options.as_blob.unwrap_or(false);
782 let guid = create_options.guid;
783 let low_32_bit_object_ids = create_options.restrict_inode_ids_to_32_bit.unwrap_or(false);
784 let volume = guard
785 .create_or_mount_volume(
786 name,
787 crypt,
788 Mode::Create { guid, low_32_bit_object_ids },
789 as_blob,
790 )
791 .await?;
792 self.serve_volume(&volume, outgoing_directory_server_end, as_blob)
793 .context("failed to serve volume")?;
794 guard.auto_unmount(volume.volume().store().store_object_id());
795 Ok(())
796 }
797
798 async fn handle_volume_requests(
799 self: &Arc<Self>,
800 name: &str,
801 mut requests: VolumeRequestStream,
802 store_id: u64,
803 ) -> Result<(), Error> {
804 while let Some(request) = requests.try_next().await? {
805 match request {
806 VolumeRequest::Check { responder, options } => {
807 async move {
808 responder.send(self.handle_check(store_id, options).await.map_err(
809 |error| {
810 error!(error:?, store_id; "Failed to check volume");
811 map_to_raw_status(error)
812 },
813 ))
814 }
815 .trace(trace_future_args!("Volume::Check"))
816 .await?;
817 }
818 VolumeRequest::Mount { responder, outgoing_directory, options } => {
819 async move {
820 responder.send(
821 self.handle_mount(name, store_id, outgoing_directory, options)
822 .await
823 .map_err(|error| {
824 error!(error:?, name, store_id; "Failed to mount volume");
825 map_to_raw_status(error)
826 }),
827 )
828 }
829 .trace(trace_future_args!("Volume::Mount"))
830 .await?;
831 }
832 VolumeRequest::SetLimit { responder, bytes } => {
833 async move {
834 responder.send(self.handle_set_limit(store_id, bytes).await.map_err(
835 |error| {
836 error!(error:?, store_id; "Failed to set volume limit");
837 map_to_raw_status(error)
838 },
839 ))
840 }
841 .trace(trace_future_args!("Volume::SetLimit"))
842 .await?;
843 }
844 VolumeRequest::GetLimit { responder } => {
845 fxfs_trace::duration!("Volume::GetLimit");
846 responder.send(Ok(self.handle_get_limit(store_id)))?
847 }
848 VolumeRequest::GetInfo { responder } => {
849 async move {
850 let result = self.handle_get_info(store_id).await.map(|guid| {
851 fidl_fuchsia_fs_startup::VolumeInfo {
852 guid: Some(guid),
853 ..Default::default()
854 }
855 });
856 match result {
857 Ok(response) => responder.send(Ok(&response)),
858 Err(error) => {
859 error!(error:?, store_id; "Failed to get volume info");
860 responder.send(Err(map_to_raw_status(error)))
861 }
862 }
863 }
864 .trace(trace_future_args!("Volume::GetInfo"))
865 .await?;
866 }
867 }
868 }
869 Ok(())
870 }
871
872 pub fn memory_pressure_config(&self) -> &MemoryPressureConfig {
873 &self.memory_pressure_config
874 }
875
876 fn is_flush_required_to_dirty(&self, byte_count: u64) -> bool {
877 let mem_pressure = self
878 .mem_monitor
879 .as_ref()
880 .map(|mem_monitor| mem_monitor.level())
881 .unwrap_or(MemoryPressureLevel::Normal);
882 if !matches!(mem_pressure, MemoryPressureLevel::Critical) {
883 return false;
884 }
885
886 let total_dirty = self.pager_dirty_bytes_count.load();
887 total_dirty + byte_count >= self.max_dirty_bytes_when_critical.load(Ordering::Relaxed)
888 }
889
890 pub fn report_pager_dirty(
895 self: Arc<Self>,
896 byte_count: u64,
897 volume: Arc<FxVolume>,
898 mark_dirty: impl FnOnce() + Send + 'static,
899 ) {
900 if !self.is_flush_required_to_dirty(byte_count) {
901 self.pager_dirty_bytes_count.fetch_add(byte_count);
902 mark_dirty();
903 } else {
904 volume.spawn(
905 async move {
906 let volumes = self.mounted_volumes.lock().await;
907
908 if self.is_flush_required_to_dirty(byte_count) {
911 debug!(
912 "Flushing all volumes. Memory pressure is critical & dirty pager bytes \
913 ({} MiB) >= limit ({} MiB)",
914 self.pager_dirty_bytes_count.load() / MEBIBYTE,
915 self.max_dirty_bytes_when_critical.load(Ordering::Relaxed) / MEBIBYTE
916 );
917
918 let flushes = FuturesUnordered::new();
919 for MountedVolume { volume, .. } in volumes.values() {
920 let vol = volume.volume().clone();
921 flushes.push(async move {
922 vol.minimize_memory().await;
923 });
924 }
925
926 flushes.collect::<()>().await;
927 }
928 self.pager_dirty_bytes_count.fetch_add(byte_count);
929 mark_dirty();
930 }
931 .trace(trace_future_args!("flush-before-mark-dirty")),
932 )
933 }
934 }
935
936 pub fn report_pager_clean(&self, byte_count: u64) {
938 let prev_dirty = self.pager_dirty_bytes_count.fetch_sub(byte_count);
939 debug_assert!(prev_dirty >= byte_count, "Underflowed dirty bytes.");
941
942 if prev_dirty < byte_count {
943 self.pager_dirty_bytes_count.store(0);
946 }
947 }
948
949 async fn handle_check(
950 self: &Arc<Self>,
951 store_id: u64,
952 options: CheckOptions,
953 ) -> Result<(), Error> {
954 let fs = self.root_volume.volume_directory().store().filesystem();
955 let crypt = if let Some(crypt) = options.crypt {
956 Some(Arc::new(RemoteCrypt::new(crypt)) as Arc<dyn Crypt>)
957 } else {
958 None
959 };
960 let result = fsck::fsck_volume(fs.as_ref(), store_id, crypt).await?;
961 info!(store_id:%; "{result:?}");
963 Ok(())
964 }
965
966 async fn handle_set_limit(self: &Arc<Self>, store_id: u64, bytes: u64) -> Result<(), Error> {
967 let store = self.root_volume.volume_directory().store();
968 let mut transaction = store.new_transaction(lock_keys![], Options::default()).await?;
969 store.filesystem().allocator().set_bytes_limit(&mut transaction, store_id, bytes)?;
970 transaction.commit().await?;
971 Ok(())
972 }
973
974 fn handle_get_limit(self: &Arc<Self>, store_id: u64) -> u64 {
975 let fs = self.root_volume.volume_directory().store().filesystem();
976 fs.allocator().get_owner_bytes_limit(store_id).unwrap_or_default()
977 }
978
979 async fn handle_get_info(self: &Arc<Self>, store_id: u64) -> Result<[u8; 16], Error> {
980 let fs = self.root_volume.volume_directory().store().filesystem();
981 let store =
982 fs.object_manager().store(store_id).ok_or_else(|| anyhow!("Store not found"))?;
983 Ok(store.guid())
984 }
985
986 async fn handle_mount(
987 self: &Arc<Self>,
988 name: &str,
989 store_id: u64,
990 outgoing_directory_server_end: ServerEnd<fio::DirectoryMarker>,
991 options: MountOptions,
992 ) -> Result<(), Error> {
993 info!(name:%, store_id:%, options:?; "Received mount request");
994 let crypt = options.crypt.map(|crypt| Arc::new(RemoteCrypt::new(crypt)) as Arc<dyn Crypt>);
995 let as_blob = options.as_blob.unwrap_or(false);
996 let mut guard = self.lock().await;
997 let volume = guard
998 .create_or_mount_volume(name, crypt, Mode::Mount, as_blob)
999 .await
1000 .context("failed to mount volume")?;
1001 self.serve_volume(&volume, outgoing_directory_server_end, as_blob)
1002 .context("failed to serve volume")?;
1003 guard.auto_unmount(volume.volume().store().store_object_id());
1004 Ok(())
1005 }
1006
1007 async fn handle_admin_requests(
1008 self: &Arc<Self>,
1009 mut stream: AdminRequestStream,
1010 store_id: u64,
1011 ) -> Result<(), Error> {
1012 if let Some(request) = stream.try_next().await.context("Reading request")? {
1014 match request {
1015 AdminRequest::Shutdown { responder } => {
1016 info!(store_id; "Received shutdown request for volume");
1017
1018 let root_store = self.root_volume.volume_directory().store();
1019 let fs = root_store.filesystem();
1020 let _guard = fs
1021 .lock_manager()
1022 .txn_lock(lock_keys![LockKey::object(
1023 root_store.store_object_id(),
1024 self.root_volume.volume_directory().object_id(),
1025 )])
1026 .await;
1027
1028 let maybe_volume = self.lock().await.unmount(store_id).await;
1029 responder
1030 .send()
1031 .unwrap_or_else(|e| warn!("Failed to send shutdown response: {}", e));
1032
1033 if let Ok(volume) = maybe_volume {
1034 volume.admin_scope().shutdown();
1037 }
1038
1039 return Ok(());
1040 }
1041 }
1042 }
1043 Ok(())
1044 }
1045
1046 pub fn set_on_mount_callback<
1051 F: Fn(&str, Option<(Arc<FxVolume>, Arc<ObjectStore>)>) + Send + Sync + 'static,
1052 >(
1053 &self,
1054 callback: F,
1055 ) {
1056 self.on_volume_added.set(Box::new(callback)).ok().unwrap();
1057 }
1058
1059 pub async fn install_volume(
1060 self: &Arc<Self>,
1061 src: &str,
1062 image_file: &str,
1063 dst: &str,
1064 ) -> Result<(), Error> {
1065 let guard = self.lock().await;
1066 info!("installing {src}/{image_file} -> {dst}");
1067 for MountedVolume { volume, .. } in guard.mounted_volumes.values() {
1068 if volume.volume().name() == src {
1069 return Err(zx::Status::ALREADY_BOUND)
1070 .with_context(|| format!("volume {src} is already mounted"));
1071 }
1072 if volume.volume().name() == dst {
1073 return Err(zx::Status::ALREADY_BOUND)
1074 .with_context(|| format!("volume {dst} is already mounted"));
1075 }
1076 }
1077 guard.volumes_directory.root_volume.install_volume(&src, &image_file, &dst).await?;
1078
1079 guard
1082 .volumes_directory
1083 .directory_node()
1084 .remove_entry(src, false)
1085 .unwrap();
1086 guard
1087 .volumes_directory
1088 .directory_node()
1089 .remove_entry(dst, false)
1090 .unwrap();
1091 let new_dst_object_id =
1092 match guard.volumes_directory.root_volume.volume_directory().lookup(dst).await? {
1093 Some((object_id, ObjectDescriptor::Volume, _)) => Ok(object_id),
1094 Some(_) => Err(FxfsError::Inconsistent),
1095 None => Err(FxfsError::NotFound),
1096 }?;
1097 self.add_directory_entry(dst, new_dst_object_id);
1098
1099 info!("install complete");
1100 Ok(())
1101 }
1102}
1103
1104#[cfg(test)]
1105pub(crate) fn serve_startup_volume_proxy(
1106 volumes_directory: &Arc<VolumesDirectory>,
1107 volume_name: &str,
1108) -> (fidl_fuchsia_fs_startup::VolumeProxy, vfs::ExecutionScope) {
1109 use vfs::ToObjectRequest;
1110 use vfs::service::ServiceLike;
1111 let scope = vfs::ExecutionScope::new();
1112 let entry = volumes_directory.directory_node().get_entry(volume_name).unwrap();
1113 let service = entry.into_any().downcast::<vfs::service::Service>().unwrap();
1114 let (proxy, server) = fidl::endpoints::create_proxy::<fidl_fuchsia_fs_startup::VolumeMarker>();
1115 service
1116 .connect(
1117 scope.clone(),
1118 Default::default(),
1119 &mut fio::Flags::PROTOCOL_SERVICE.to_object_request(server),
1120 )
1121 .unwrap();
1122 (proxy, scope)
1123}
1124
1125struct PagerDirtyByteCount(AtomicU64);
1126
1127impl PagerDirtyByteCount {
1128 pub fn new() -> Self {
1129 Self(AtomicU64::new(0))
1130 }
1131
1132 pub fn fetch_add(&self, value: u64) -> u64 {
1133 let prev = self.0.fetch_add(value, Ordering::Relaxed);
1134 fxfs_trace::counter!("dirty-bytes", 0, "total" => prev.saturating_add(value));
1135 prev
1136 }
1137
1138 pub fn fetch_sub(&self, value: u64) -> u64 {
1139 let prev = self.0.fetch_sub(value, Ordering::Relaxed);
1140 fxfs_trace::counter!("dirty-bytes", 0, "total" => prev.saturating_sub(value));
1141 prev
1142 }
1143
1144 pub fn load(&self) -> u64 {
1145 self.0.load(Ordering::Relaxed)
1146 }
1147
1148 pub fn store(&self, value: u64) {
1149 self.0.store(value, Ordering::Relaxed);
1150 fxfs_trace::counter!("dirty-bytes", 0, "total" => value);
1151 }
1152}
1153
1154#[cfg(test)]
1155mod tests {
1156 use super::{Mode, serve_startup_volume_proxy};
1157 use crate::fuchsia::RemoteCrypt;
1158 use crate::fuchsia::memory_pressure::MemoryPressureLevel;
1159 use crate::fuchsia::testing::{self, TestFixture, open_dir_checked, open_file_checked};
1160 use crate::fuchsia::volume::MemoryPressureConfig;
1161 use crate::fuchsia::volumes_directory::VolumesDirectory;
1162 use crate::testing::TestFixtureOptions;
1163 use fidl::endpoints::{DiscoverableProtocolMarker, create_proxy, create_request_stream};
1164 use fidl_fuchsia_fs::AdminMarker;
1165 use fidl_fuchsia_fs_startup::{MountOptions, VolumeProxy};
1166 use fidl_fuchsia_fxfs::{CryptRequest, FxfsKey, KeyPurpose, WrappedKey};
1167 use fidl_fuchsia_io as fio;
1168 use fuchsia_async as fasync;
1169 use fuchsia_component_client::connect_to_protocol_at_dir_svc;
1170 use fuchsia_fs::file;
1171 use futures::{TryStreamExt, join};
1172 use fxfs::errors::FxfsError;
1173 use fxfs::filesystem::FxFilesystem;
1174 use fxfs::fsck::{FsckOptions, fsck, fsck_volume_with_options, fsck_with_options};
1175 use fxfs::lock_keys;
1176 use fxfs::object_handle::ObjectHandle;
1177 use fxfs::object_store::allocator::Allocator;
1178 use fxfs::object_store::transaction::{LockKey, Options};
1179 use fxfs::object_store::volume::root_volume;
1180 use fxfs_crypto::Crypt;
1181 use fxfs_insecure_crypto::new_insecure_crypt;
1182 use refaults_vmo::PageRefaultCounter;
1183 use std::sync::atomic::Ordering;
1184 use std::sync::{Arc, Weak};
1185 use std::time::Duration;
1186 use storage_device::DeviceHolder;
1187 use storage_device::fake_device::FakeDevice;
1188 use vfs::execution_scope::ExecutionScope;
1189 use vfs::temp_clone::{TempClonable, unblock};
1190 use zx::Status;
1191 async fn write_image_to_file(image: DeviceHolder, file: fio::FileProxy) {
1192 file.resize(image.size()).await.unwrap().expect("resize failed");
1193 let vmo = TempClonable::new(
1194 file.get_backing_memory(fio::VmoFlags::SHARED_BUFFER | fio::VmoFlags::WRITE)
1195 .await
1196 .unwrap()
1197 .expect("get backing memory failed"),
1198 );
1199
1200 const CHUNK_READ_SIZE: usize = 131_072; let mut buff = image.allocate_buffer(CHUNK_READ_SIZE).await;
1202 let total = image.size();
1203 let mut offset = 0;
1204 while offset < total {
1205 let amount = std::cmp::min(total - offset, CHUNK_READ_SIZE as u64);
1206 image.read(offset, buff.as_mut()).await.expect("image read failed");
1207 {
1208 let vmo = vmo.temp_clone();
1211 let data = buff.as_slice()[0..amount as usize].to_vec();
1212 let offset = offset;
1213 unblock(move || vmo.write(&data, offset)).await.expect("vmo write failed");
1214 }
1215 offset += amount;
1216 }
1217 assert_eq!(offset, total);
1218 file.sync().await.unwrap().expect("sync failed");
1219 file.close().await.unwrap().expect("close failed");
1220 }
1221
1222 #[fuchsia::test]
1223 async fn test_volume_creation() {
1224 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1225 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1226 let blob_resupplied_count =
1227 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1228 let volumes_directory = VolumesDirectory::new(
1229 root_volume(filesystem.clone()).await.unwrap(),
1230 Weak::new(),
1231 None,
1232 blob_resupplied_count,
1233 MemoryPressureConfig::default(),
1234 )
1235 .await
1236 .unwrap();
1237
1238 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1239 {
1240 let vol = volumes_directory
1241 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1242 .await
1243 .expect("create encrypted volume failed");
1244 vol.volume().store().store_object_id()
1245 };
1246
1247 volumes_directory.terminate().await;
1248 std::mem::drop(volumes_directory);
1249 filesystem.close().await.expect("close filesystem failed");
1250 let device = filesystem.take_device().await;
1251 device.reopen(false);
1252 let filesystem = FxFilesystem::open(device).await.unwrap();
1253 let blob_resupplied_count =
1254 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1255 let volumes_directory = VolumesDirectory::new(
1256 root_volume(filesystem.clone()).await.unwrap(),
1257 Weak::new(),
1258 None,
1259 blob_resupplied_count,
1260 MemoryPressureConfig::default(),
1261 )
1262 .await
1263 .unwrap();
1264
1265 let error = volumes_directory
1266 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1267 .await
1268 .err()
1269 .expect("Creating existing encrypted volume should fail");
1270 assert!(FxfsError::AlreadyExists.matches(&error));
1271 }
1272
1273 #[fuchsia::test]
1274 async fn test_dirty_pages_accumulate_in_parent() {
1275 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1276 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1277 let blob_resupplied_count =
1278 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1279 let volumes_directory = VolumesDirectory::new(
1280 root_volume(filesystem.clone()).await.unwrap(),
1281 Weak::new(),
1282 None,
1283 blob_resupplied_count,
1284 MemoryPressureConfig::default(),
1285 )
1286 .await
1287 .unwrap();
1288
1289 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1290 let vol = volumes_directory
1291 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1292 .await
1293 .expect("create encrypted volume failed");
1294 let old_dirty = volumes_directory.pager_dirty_bytes_count.load();
1295
1296 let new_dirty = {
1297 let (root, server_end) = create_proxy::<fio::DirectoryMarker>();
1298 vol.root().clone().serve(fio::PERM_READABLE | fio::PERM_WRITABLE, server_end);
1299 let f = open_file_checked(
1300 &root,
1301 "foo",
1302 fio::Flags::FLAG_MAYBE_CREATE
1303 | fio::PERM_READABLE
1304 | fio::PERM_WRITABLE
1305 | fio::Flags::PROTOCOL_FILE,
1306 &Default::default(),
1307 )
1308 .await;
1309 let buf = vec![0xaa as u8; 8192];
1310 file::write(&f, buf.as_slice()).await.expect("Write");
1311 volumes_directory.pager_dirty_bytes_count.load()
1314 };
1315 assert_ne!(old_dirty, new_dirty);
1316
1317 volumes_directory.terminate().await;
1318 std::mem::drop(volumes_directory);
1319 filesystem.close().await.expect("close filesystem failed");
1320 }
1321
1322 #[fuchsia::test]
1323 async fn test_volume_reopen() {
1324 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1325 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1326 let blob_resupplied_count =
1327 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1328 let volumes_directory = VolumesDirectory::new(
1329 root_volume(filesystem.clone()).await.unwrap(),
1330 Weak::new(),
1331 None,
1332 blob_resupplied_count,
1333 MemoryPressureConfig::default(),
1334 )
1335 .await
1336 .unwrap();
1337
1338 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1339 let volume_id = {
1340 let vol = volumes_directory
1341 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1342 .await
1343 .expect("create encrypted volume failed");
1344 vol.volume().store().store_object_id()
1345 };
1346
1347 volumes_directory.terminate().await;
1348 std::mem::drop(volumes_directory);
1349 filesystem.close().await.expect("close filesystem failed");
1350 let device = filesystem.take_device().await;
1351 device.reopen(false);
1352 let filesystem = FxFilesystem::open(device).await.unwrap();
1353 let blob_resupplied_count =
1354 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1355 let volumes_directory = VolumesDirectory::new(
1356 root_volume(filesystem.clone()).await.unwrap(),
1357 Weak::new(),
1358 None,
1359 blob_resupplied_count,
1360 MemoryPressureConfig::default(),
1361 )
1362 .await
1363 .unwrap();
1364
1365 {
1366 let vol = volumes_directory
1367 .mount_volume("encrypted", Some(crypt.clone()), false)
1368 .await
1369 .expect("open existing encrypted volume failed");
1370 assert_eq!(vol.volume().store().store_object_id(), volume_id);
1371 }
1372
1373 volumes_directory.terminate().await;
1374 std::mem::drop(volumes_directory);
1375 filesystem.close().await.expect("close filesystem failed");
1376 }
1377
1378 #[fuchsia::test]
1379 async fn test_volume_creation_unencrypted() {
1380 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1381 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1382 let blob_resupplied_count =
1383 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1384 let volumes_directory = VolumesDirectory::new(
1385 root_volume(filesystem.clone()).await.unwrap(),
1386 Weak::new(),
1387 None,
1388 blob_resupplied_count,
1389 MemoryPressureConfig::default(),
1390 )
1391 .await
1392 .unwrap();
1393
1394 {
1395 let vol = volumes_directory
1396 .create_and_mount_volume("unencrypted", None, false, None)
1397 .await
1398 .expect("create unencrypted volume failed");
1399 vol.volume().store().store_object_id()
1400 };
1401
1402 volumes_directory.terminate().await;
1403 std::mem::drop(volumes_directory);
1404 filesystem.close().await.expect("close filesystem failed");
1405 let device = filesystem.take_device().await;
1406 device.reopen(false);
1407 let filesystem = FxFilesystem::open(device).await.unwrap();
1408 let blob_resupplied_count =
1409 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1410 let volumes_directory = VolumesDirectory::new(
1411 root_volume(filesystem.clone()).await.unwrap(),
1412 Weak::new(),
1413 None,
1414 blob_resupplied_count,
1415 MemoryPressureConfig::default(),
1416 )
1417 .await
1418 .unwrap();
1419
1420 let error = volumes_directory
1421 .create_and_mount_volume("unencrypted", None, false, None)
1422 .await
1423 .err()
1424 .expect("Creating existing unencrypted volume should fail");
1425 assert!(FxfsError::AlreadyExists.matches(&error));
1426
1427 volumes_directory.terminate().await;
1428 std::mem::drop(volumes_directory);
1429 filesystem.close().await.expect("close filesystem failed");
1430 }
1431
1432 #[fuchsia::test]
1433 async fn test_volume_reopen_unencrypted() {
1434 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1435 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1436 let blob_resupplied_count =
1437 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1438 let volumes_directory = VolumesDirectory::new(
1439 root_volume(filesystem.clone()).await.unwrap(),
1440 Weak::new(),
1441 None,
1442 blob_resupplied_count,
1443 MemoryPressureConfig::default(),
1444 )
1445 .await
1446 .unwrap();
1447
1448 let volume_id = {
1449 let vol = volumes_directory
1450 .create_and_mount_volume("unencrypted", None, false, None)
1451 .await
1452 .expect("create unencrypted volume failed");
1453 vol.volume().store().store_object_id()
1454 };
1455
1456 volumes_directory.terminate().await;
1457 std::mem::drop(volumes_directory);
1458 filesystem.close().await.expect("close filesystem failed");
1459 let device = filesystem.take_device().await;
1460 device.reopen(false);
1461 let filesystem = FxFilesystem::open(device).await.unwrap();
1462 let blob_resupplied_count =
1463 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1464 let volumes_directory = VolumesDirectory::new(
1465 root_volume(filesystem.clone()).await.unwrap(),
1466 Weak::new(),
1467 None,
1468 blob_resupplied_count,
1469 MemoryPressureConfig::default(),
1470 )
1471 .await
1472 .unwrap();
1473
1474 {
1475 let vol = volumes_directory
1476 .mount_volume("unencrypted", None, false)
1477 .await
1478 .expect("open existing unencrypted volume failed");
1479 assert_eq!(vol.volume().store().store_object_id(), volume_id);
1480 }
1481
1482 volumes_directory.terminate().await;
1483 std::mem::drop(volumes_directory);
1484 filesystem.close().await.expect("close filesystem failed");
1485 }
1486
1487 #[fuchsia::test]
1488 async fn test_volume_enumeration() {
1489 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1490 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1491 let blob_resupplied_count =
1492 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1493 let volumes_directory = VolumesDirectory::new(
1494 root_volume(filesystem.clone()).await.unwrap(),
1495 Weak::new(),
1496 None,
1497 blob_resupplied_count,
1498 MemoryPressureConfig::default(),
1499 )
1500 .await
1501 .unwrap();
1502
1503 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1505 {
1506 volumes_directory
1507 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1508 .await
1509 .expect("create encrypted volume failed");
1510 };
1511 {
1513 volumes_directory
1514 .create_and_mount_volume("unencrypted", None, false, None)
1515 .await
1516 .expect("create unencrypted volume failed");
1517 };
1518
1519 volumes_directory.terminate().await;
1521 std::mem::drop(volumes_directory);
1522 filesystem.close().await.expect("close filesystem failed");
1523 let device = filesystem.take_device().await;
1524 device.reopen(false);
1525 let filesystem = FxFilesystem::open(device).await.unwrap();
1526 let blob_resupplied_count =
1527 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1528 let volumes_directory = VolumesDirectory::new(
1529 root_volume(filesystem.clone()).await.unwrap(),
1530 Weak::new(),
1531 None,
1532 blob_resupplied_count,
1533 MemoryPressureConfig::default(),
1534 )
1535 .await
1536 .unwrap();
1537
1538 let readdir = |dir: Arc<fio::DirectoryProxy>| async move {
1539 let status = dir.rewind().await.expect("FIDL call failed");
1540 Status::ok(status).expect("rewind failed");
1541 let (status, buf) = dir.read_dirents(fio::MAX_BUF).await.expect("FIDL call failed");
1542 Status::ok(status).expect("read_dirents failed");
1543 let mut entries = vec![];
1544 for res in fuchsia_fs::directory::parse_dir_entries(&buf) {
1545 entries.push(res.expect("Failed to parse entry").name);
1546 }
1547 entries
1548 };
1549
1550 let dir_proxy = Arc::new(vfs::directory::serve_read_only(
1551 volumes_directory.directory_node().clone(),
1552 ExecutionScope::new(),
1553 ));
1554 let entries = readdir(dir_proxy.clone()).await;
1555 assert_eq!(entries, [".", "encrypted", "unencrypted"]);
1556
1557 let _vol = volumes_directory
1558 .mount_volume("encrypted", Some(crypt.clone()), false)
1559 .await
1560 .expect("Open encrypted volume failed");
1561
1562 let entries = readdir(dir_proxy).await;
1564 assert_eq!(entries, [".", "encrypted", "unencrypted"]);
1565
1566 volumes_directory.terminate().await;
1567 std::mem::drop(volumes_directory);
1568 filesystem.close().await.expect("close filesystem failed");
1569 }
1570
1571 #[fuchsia::test]
1572 async fn test_get_info() {
1573 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1574 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1575 let root_volume = root_volume(filesystem.clone()).await.unwrap();
1576 let blob_resupplied_count =
1577 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1578 let volumes_directory = VolumesDirectory::new(
1579 root_volume,
1580 Weak::new(),
1581 None,
1582 blob_resupplied_count,
1583 MemoryPressureConfig::default(),
1584 )
1585 .await
1586 .unwrap();
1587
1588 let vol = volumes_directory
1589 .create_and_mount_volume("vol", None, false, None)
1590 .await
1591 .expect("create_and_mount_volume failed");
1592 let guid = vol.volume().store().guid();
1593
1594 let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, "vol");
1595
1596 let info: fidl_fuchsia_fs_startup::VolumeInfo = volume_proxy
1597 .get_info()
1598 .await
1599 .expect("get_info failed")
1600 .expect("get_info returned error");
1601 assert_eq!(info.guid, Some(guid));
1602
1603 volumes_directory.terminate().await;
1604 }
1605
1606 #[fuchsia::test]
1607 async fn test_deleted_encrypted_volume_while_mounted() {
1608 const VOLUME_NAME: &str = "encrypted";
1609
1610 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1611 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1612 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1613 let blob_resupplied_count =
1614 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1615 let volumes_directory = VolumesDirectory::new(
1616 root_volume(filesystem.clone()).await.unwrap(),
1617 Weak::new(),
1618 None,
1619 blob_resupplied_count,
1620 MemoryPressureConfig::default(),
1621 )
1622 .await
1623 .unwrap();
1624 volumes_directory
1625 .create_and_mount_volume(VOLUME_NAME, Some(crypt.clone()), false, None)
1626 .await
1627 .expect("create encrypted volume failed");
1628 assert!(
1630 FxfsError::AlreadyBound.matches(
1631 &volumes_directory
1632 .remove_volume(VOLUME_NAME)
1633 .await
1634 .err()
1635 .expect("Deleting volume should fail")
1636 )
1637 );
1638 volumes_directory.terminate().await;
1639 std::mem::drop(volumes_directory);
1640 filesystem.close().await.expect("close filesystem failed");
1641 }
1642
1643 #[fuchsia::test]
1644 async fn test_mount_volume_using_volume_protocol() {
1645 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1646 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1647 let blob_resupplied_count =
1648 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1649 let volumes_directory = VolumesDirectory::new(
1650 root_volume(filesystem.clone()).await.unwrap(),
1651 Weak::new(),
1652 None,
1653 blob_resupplied_count,
1654 MemoryPressureConfig::default(),
1655 )
1656 .await
1657 .unwrap();
1658
1659 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1660 let store_id = {
1661 let vol = volumes_directory
1662 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1663 .await
1664 .expect("create encrypted volume failed");
1665 vol.volume().store().store_object_id()
1666 };
1667 volumes_directory.lock().await.unmount(store_id).await.expect("unmount failed");
1668
1669 let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, "encrypted");
1670
1671 let (dir_proxy, dir_server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1672
1673 let crypt_service = fxfs_crypt::CryptService::new();
1674 crypt_service
1675 .add_wrapping_key(0, fxfs_insecure_crypto::DATA_KEY.to_vec())
1676 .expect("add_wrapping_key failed");
1677 crypt_service
1678 .add_wrapping_key(1, fxfs_insecure_crypto::METADATA_KEY.to_vec())
1679 .expect("add_wrapping_key failed");
1680 crypt_service.set_active_key(KeyPurpose::Data, 0).expect("set_active_key failed");
1681 crypt_service.set_active_key(KeyPurpose::Metadata, 1).expect("set_active_key failed");
1682 let (client1, stream1) = create_request_stream();
1683 let (client2, stream2) = create_request_stream();
1684
1685 join!(
1686 async {
1687 volume_proxy
1688 .mount(
1689 dir_server_end,
1690 MountOptions { crypt: Some(client1), ..MountOptions::default() },
1691 )
1692 .await
1693 .expect("mount (fidl) failed")
1694 .expect("mount failed");
1695
1696 open_file_checked(
1697 &dir_proxy,
1698 "root/test",
1699 fio::PERM_READABLE | fio::PERM_WRITABLE | fio::Flags::FLAG_MAYBE_CREATE,
1700 &Default::default(),
1701 )
1702 .await;
1703
1704 let (_dir_proxy, dir_server_end) =
1706 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1707
1708 assert_eq!(
1709 Status::from_raw(
1710 volume_proxy
1711 .mount(
1712 dir_server_end,
1713 MountOptions { crypt: Some(client2), ..MountOptions::default() },
1714 )
1715 .await
1716 .expect("mount (fidl) failed")
1717 .expect_err("mount succeeded")
1718 ),
1719 Status::ALREADY_BOUND
1720 );
1721
1722 std::mem::drop(dir_proxy);
1723
1724 let mut count = 0;
1726 loop {
1727 if volumes_directory.mounted_volumes.lock().await.is_empty() {
1728 break;
1729 }
1730 count += 1;
1731 assert!(count <= 100);
1732 fasync::Timer::new(Duration::from_millis(100)).await;
1733 }
1734 },
1735 async {
1736 crypt_service
1737 .handle_request(fxfs_crypt::Services::Crypt(stream1))
1738 .await
1739 .expect("handle_request failed");
1740 crypt_service
1741 .handle_request(fxfs_crypt::Services::Crypt(stream2))
1742 .await
1743 .expect("handle_request failed");
1744 }
1745 );
1746 volumes_directory.terminate().await;
1752 }
1753
1754 #[fuchsia::test]
1755 #[ignore] async fn test_volume_dir_races() {
1758 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1759 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1760 let blob_resupplied_count =
1761 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1762 let volumes_directory = VolumesDirectory::new(
1763 root_volume(filesystem.clone()).await.unwrap(),
1764 Weak::new(),
1765 None,
1766 blob_resupplied_count,
1767 MemoryPressureConfig::default(),
1768 )
1769 .await
1770 .unwrap();
1771
1772 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1773 let store_id = {
1774 let vol = volumes_directory
1775 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1776 .await
1777 .expect("create encrypted volume failed");
1778 vol.volume().store().store_object_id()
1779 };
1780 volumes_directory.lock().await.unmount(store_id).await.expect("unmount failed");
1781
1782 let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, "encrypted");
1783
1784 let crypt_service = Arc::new(fxfs_crypt::CryptService::new());
1785 crypt_service
1786 .add_wrapping_key(0, fxfs_insecure_crypto::DATA_KEY.to_vec())
1787 .expect("add_wrapping_key failed");
1788 crypt_service
1789 .add_wrapping_key(1, fxfs_insecure_crypto::METADATA_KEY.to_vec())
1790 .expect("add_wrapping_key failed");
1791 crypt_service.set_active_key(KeyPurpose::Data, 0).expect("set_active_key failed");
1792 crypt_service.set_active_key(KeyPurpose::Metadata, 1).expect("set_active_key failed");
1793 let (client1, stream1) = create_request_stream();
1794 let (client2, stream2) = create_request_stream();
1795 let crypt_service_clone = crypt_service.clone();
1796 let crypt_task1 = fasync::Task::spawn(async move {
1797 crypt_service_clone
1798 .handle_request(fxfs_crypt::Services::Crypt(stream1))
1799 .await
1800 .expect("handle_request failed");
1801 });
1802 let crypt_task2 = fasync::Task::spawn(async move {
1803 crypt_service
1804 .handle_request(fxfs_crypt::Services::Crypt(stream2))
1805 .await
1806 .expect("handle_request failed");
1807 });
1808
1809 join!(
1813 async {
1814 let (_dir_proxy, dir_server_end) =
1815 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1816 if let Err(status) = volume_proxy
1817 .mount(
1818 dir_server_end,
1819 MountOptions { crypt: Some(client1), ..MountOptions::default() },
1820 )
1821 .await
1822 .expect("mount (fidl) failed")
1823 {
1824 let status = Status::from_raw(status);
1825 if status != Status::NOT_FOUND && status != Status::ALREADY_BOUND {
1826 assert!(false, "Unexpected status {:}", status);
1827 }
1828 }
1829 },
1830 async {
1831 let (_dir_proxy, dir_server_end) =
1832 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1833 if let Err(status) = volume_proxy
1834 .mount(
1835 dir_server_end,
1836 MountOptions { crypt: Some(client2), ..MountOptions::default() },
1837 )
1838 .await
1839 .expect("mount (fidl) failed")
1840 {
1841 let status = Status::from_raw(status);
1842 if status != Status::NOT_FOUND && status != Status::ALREADY_BOUND {
1843 assert!(false, "Unexpected status {:}", status);
1844 }
1845 }
1846 },
1847 async {
1848 let volumes_directory = volumes_directory.clone();
1849 let wait_time = rand::random_range(0..5);
1850 fasync::Timer::new(Duration::from_millis(wait_time)).await;
1851 if let Err(err) = volumes_directory.remove_volume("encrypted").await {
1852 assert!(
1853 FxfsError::NotFound.matches(&err) || FxfsError::AlreadyBound.matches(&err),
1854 "Unexpected error {:?}",
1855 err
1856 );
1857 }
1858 },
1859 async {
1860 let volumes_directory = volumes_directory.clone();
1861 let wait_time = rand::random_range(0..5);
1862 fasync::Timer::new(Duration::from_millis(wait_time)).await;
1863 if let Err(err) = volumes_directory.remove_volume("encrypted").await {
1864 assert!(
1865 FxfsError::NotFound.matches(&err) || FxfsError::AlreadyBound.matches(&err),
1866 "Unexpected error {:?}",
1867 err
1868 );
1869 }
1870 },
1871 async {
1872 let volumes_directory = volumes_directory.clone();
1873 let wait_time = rand::random_range(0..5);
1874 fasync::Timer::new(Duration::from_millis(wait_time)).await;
1875 let mut guard = volumes_directory.lock().await;
1876 match guard
1877 .create_or_mount_volume(
1878 "encrypted",
1879 Some(crypt.clone()),
1880 Mode::Create { guid: None, low_32_bit_object_ids: false },
1881 false,
1882 )
1883 .await
1884 {
1885 Ok(vol) => {
1886 let store_id = vol.volume().store().store_object_id();
1887 std::mem::drop(vol);
1888 guard.unmount(store_id).await.expect("unmount failed");
1889 }
1890 Err(err) => {
1891 assert!(
1892 FxfsError::AlreadyExists.matches(&err)
1893 || FxfsError::AlreadyBound.matches(&err),
1894 "Unexpected error {:?}",
1895 err
1896 );
1897 }
1898 }
1899 }
1900 );
1901 std::mem::drop(crypt_task1);
1902 std::mem::drop(crypt_task2);
1903 volumes_directory.terminate().await;
1909 }
1910
1911 #[fuchsia::test]
1912 async fn test_shutdown_volume() {
1913 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
1914 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1915 let blob_resupplied_count =
1916 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1917 let volumes_directory = VolumesDirectory::new(
1918 root_volume(filesystem.clone()).await.unwrap(),
1919 Weak::new(),
1920 None,
1921 blob_resupplied_count,
1922 MemoryPressureConfig::default(),
1923 )
1924 .await
1925 .unwrap();
1926
1927 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
1928 let vol = volumes_directory
1929 .create_and_mount_volume("encrypted", Some(crypt.clone()), false, None)
1930 .await
1931 .expect("create encrypted volume failed");
1932
1933 let (dir_proxy, dir_server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
1934
1935 volumes_directory.serve_volume(&vol, dir_server_end, false).expect("serve_volume failed");
1936
1937 let admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(&dir_proxy)
1938 .expect("Unable to connect to admin service");
1939
1940 admin_proxy.shutdown().await.expect("shutdown failed");
1941
1942 assert!(volumes_directory.mounted_volumes.lock().await.is_empty());
1943 }
1944
1945 #[fuchsia::test]
1946 async fn test_byte_limit_persistence() {
1947 const BYTES_LIMIT_1: u64 = 123456;
1948 const BYTES_LIMIT_2: u64 = 456789;
1949 const VOLUME_NAME: &str = "A";
1950 let mut device = DeviceHolder::new(FakeDevice::new(8192, 512));
1951 {
1952 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
1953 let blob_resupplied_count =
1954 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1955 let volumes_directory = VolumesDirectory::new(
1956 root_volume(filesystem.clone()).await.unwrap(),
1957 Weak::new(),
1958 None,
1959 blob_resupplied_count,
1960 MemoryPressureConfig::default(),
1961 )
1962 .await
1963 .unwrap();
1964
1965 volumes_directory
1966 .create_and_mount_volume(VOLUME_NAME, None, false, None)
1967 .await
1968 .expect("create unencrypted volume failed");
1969
1970 let (volume_proxy, _scope) =
1971 serve_startup_volume_proxy(&volumes_directory, VOLUME_NAME);
1972
1973 volume_proxy.set_limit(BYTES_LIMIT_1).await.unwrap().expect("To set limits");
1974 {
1975 let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
1976 assert_eq!(limits.len(), 1);
1977 assert_eq!(limits[0].1, BYTES_LIMIT_1);
1978 }
1979
1980 volume_proxy.set_limit(BYTES_LIMIT_2).await.unwrap().expect("To set limits");
1981 {
1982 let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
1983 assert_eq!(limits.len(), 1);
1984 assert_eq!(limits[0].1, BYTES_LIMIT_2);
1985 }
1986 std::mem::drop(volume_proxy);
1987 volumes_directory.terminate().await;
1988 std::mem::drop(volumes_directory);
1989 filesystem.close().await.expect("close filesystem failed");
1990 device = filesystem.take_device().await;
1991 }
1992 device.ensure_unique();
1993 device.reopen(false);
1994 {
1995 let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
1996 fsck(filesystem.clone()).await.expect("Fsck");
1997 let blob_resupplied_count =
1998 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
1999 let volumes_directory = VolumesDirectory::new(
2000 root_volume(filesystem.clone()).await.unwrap(),
2001 Weak::new(),
2002 None,
2003 blob_resupplied_count,
2004 MemoryPressureConfig::default(),
2005 )
2006 .await
2007 .unwrap();
2008 {
2009 let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
2010 assert_eq!(limits.len(), 1);
2011 assert_eq!(limits[0].1, BYTES_LIMIT_2);
2012 }
2013 volumes_directory.remove_volume(VOLUME_NAME).await.expect("Volume deletion failed");
2014 {
2015 let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
2016 assert_eq!(limits.len(), 0);
2017 }
2018 volumes_directory.terminate().await;
2019 std::mem::drop(volumes_directory);
2020 filesystem.close().await.expect("close filesystem failed");
2021 device = filesystem.take_device().await;
2022 }
2023 device.ensure_unique();
2024 device.reopen(false);
2025 let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
2026 fsck(filesystem.clone()).await.expect("Fsck");
2027 let limits = (filesystem.allocator() as Arc<Allocator>).owner_byte_limits();
2028 assert_eq!(limits.len(), 0);
2029 }
2030
2031 struct VolumeInfo {
2032 _scope: vfs::ExecutionScope,
2033 volume_proxy: VolumeProxy,
2034 file_proxy: fio::FileProxy,
2035 }
2036
2037 impl VolumeInfo {
2038 async fn new(volumes_directory: &Arc<VolumesDirectory>, name: &'static str) -> Self {
2039 let volume = volumes_directory
2040 .create_and_mount_volume(name, None, false, None)
2041 .await
2042 .expect("create unencrypted volume failed");
2043
2044 let (volume_dir_proxy, dir_server_end) =
2045 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2046 volumes_directory
2047 .serve_volume(&volume, dir_server_end, false)
2048 .expect("serve_volume failed");
2049
2050 let (volume_proxy, _scope) = serve_startup_volume_proxy(&volumes_directory, name);
2051
2052 let (root_proxy, root_server_end) =
2053 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2054 volume_dir_proxy
2055 .open(
2056 "root",
2057 fio::PERM_READABLE | fio::PERM_WRITABLE | fio::Flags::PROTOCOL_DIRECTORY,
2058 &Default::default(),
2059 root_server_end.into_channel(),
2060 )
2061 .expect("Failed to open volume root");
2062
2063 let file_proxy = open_file_checked(
2064 &root_proxy,
2065 "foo",
2066 fio::Flags::FLAG_MAYBE_CREATE
2067 | fio::PERM_READABLE
2068 | fio::PERM_WRITABLE
2069 | fio::Flags::PROTOCOL_FILE,
2070 &Default::default(),
2071 )
2072 .await;
2073 VolumeInfo { _scope, volume_proxy, file_proxy }
2074 }
2075 }
2076
2077 #[fuchsia::test]
2078 async fn test_limit_bytes() {
2079 const BYTES_LIMIT: u64 = 262_144; const BLOCK_SIZE: usize = 8192; let device = DeviceHolder::new(FakeDevice::new(BLOCK_SIZE.try_into().unwrap(), 512));
2082 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2083 let blob_resupplied_count =
2084 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2085 let volumes_directory = VolumesDirectory::new(
2086 root_volume(filesystem.clone()).await.unwrap(),
2087 Weak::new(),
2088 None,
2089 blob_resupplied_count,
2090 MemoryPressureConfig::default(),
2091 )
2092 .await
2093 .unwrap();
2094
2095 let vol = VolumeInfo::new(&volumes_directory, "foo").await;
2096 let old_info = {
2097 let (status, info) = vol.file_proxy.query_filesystem().await.expect("Getting fs info");
2098 assert_eq!(status, zx::Status::OK.into_raw());
2099 let info = info.unwrap();
2100 assert!(info.total_bytes > BYTES_LIMIT);
2102 info
2103 };
2104
2105 vol.volume_proxy.set_limit(BYTES_LIMIT).await.unwrap().expect("To set limits");
2106 {
2107 let (status, info) = vol.file_proxy.query_filesystem().await.expect("Getting fs info");
2108 assert!(status == zx::Status::OK.into_raw());
2109 let new_info = info.unwrap();
2110 assert_eq!(new_info.total_bytes, BYTES_LIMIT);
2111 assert!(new_info.used_bytes < old_info.used_bytes);
2114 }
2115
2116 let zeros = vec![0u8; BLOCK_SIZE];
2117 assert_eq!(
2119 <u64 as TryInto<usize>>::try_into(
2120 vol.file_proxy
2121 .write(&zeros)
2122 .await
2123 .expect("Failed Write message")
2124 .expect("Failed write")
2125 )
2126 .unwrap(),
2127 BLOCK_SIZE
2128 );
2129 for _ in (BLOCK_SIZE..BYTES_LIMIT as usize).step_by(BLOCK_SIZE) {
2131 match vol.file_proxy.write(&zeros).await.expect("Failed Write message") {
2132 Err(_) => break,
2133 Ok(b) if b < BLOCK_SIZE.try_into().unwrap() => break,
2134 _ => (),
2135 };
2136 }
2137
2138 assert_eq!(
2140 vol.file_proxy
2141 .write(&zeros)
2142 .await
2143 .expect("Failed write message")
2144 .expect_err("Write should have been limited"),
2145 Status::NO_SPACE.into_raw()
2146 );
2147
2148 vol.volume_proxy.set_limit(BYTES_LIMIT * 2).await.unwrap().expect("To set limits");
2150 assert_eq!(
2151 <u64 as TryInto<usize>>::try_into(
2152 vol.file_proxy
2153 .write(&zeros)
2154 .await
2155 .expect("Failed Write message")
2156 .expect("Failed write")
2157 )
2158 .unwrap(),
2159 BLOCK_SIZE
2160 );
2161
2162 vol.file_proxy.close().await.unwrap().expect("Failed to close file");
2163 volumes_directory.terminate().await;
2164 std::mem::drop(volumes_directory);
2165 filesystem.close().await.expect("close filesystem failed");
2166 }
2167
2168 #[fuchsia::test]
2169 async fn test_limit_bytes_two_hit_device_limit() {
2170 const BYTES_LIMIT: u64 = 3_145_728; const BLOCK_SIZE: usize = 8192; const BLOCK_COUNT: u32 = 512;
2173 let device =
2174 DeviceHolder::new(FakeDevice::new(BLOCK_SIZE.try_into().unwrap(), BLOCK_COUNT));
2175 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2176 let blob_resupplied_count =
2177 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2178 let volumes_directory = VolumesDirectory::new(
2179 root_volume(filesystem.clone()).await.unwrap(),
2180 Weak::new(),
2181 None,
2182 blob_resupplied_count,
2183 MemoryPressureConfig::default(),
2184 )
2185 .await
2186 .unwrap();
2187
2188 let a = VolumeInfo::new(&volumes_directory, "foo").await;
2189 let b = VolumeInfo::new(&volumes_directory, "bar").await;
2190 a.volume_proxy.set_limit(BYTES_LIMIT).await.unwrap().expect("To set limits");
2191 b.volume_proxy.set_limit(BYTES_LIMIT).await.unwrap().expect("To set limits");
2192 let mut a_written: u64 = 0;
2193 let mut b_written: u64 = 0;
2194
2195 let zeros = vec![0u8; BLOCK_SIZE];
2197
2198 assert_eq!(
2200 <u64 as TryInto<usize>>::try_into(
2201 a.file_proxy
2202 .write(&zeros)
2203 .await
2204 .expect("Failed Write message")
2205 .expect("Failed write")
2206 )
2207 .unwrap(),
2208 BLOCK_SIZE
2209 );
2210 a_written += BLOCK_SIZE as u64;
2211 assert_eq!(
2212 <u64 as TryInto<usize>>::try_into(
2213 b.file_proxy
2214 .write(&zeros)
2215 .await
2216 .expect("Failed Write message")
2217 .expect("Failed write")
2218 )
2219 .unwrap(),
2220 BLOCK_SIZE
2221 );
2222 b_written += BLOCK_SIZE as u64;
2223
2224 for _ in (BLOCK_SIZE..BYTES_LIMIT as usize).step_by(BLOCK_SIZE) {
2226 match a.file_proxy.write(&zeros).await.expect("Failed Write message") {
2227 Err(_) => break,
2228 Ok(bytes) => {
2229 a_written += bytes;
2230 if bytes < BLOCK_SIZE.try_into().unwrap() {
2231 break;
2232 }
2233 }
2234 };
2235 }
2236 assert_eq!(
2238 a.file_proxy
2239 .write(&zeros)
2240 .await
2241 .expect("Failed write message")
2242 .expect_err("Write should have been limited"),
2243 Status::NO_SPACE.into_raw()
2244 );
2245
2246 for _ in (BLOCK_SIZE..BYTES_LIMIT as usize).step_by(BLOCK_SIZE) {
2249 match b.file_proxy.write(&zeros).await.expect("Failed Write message") {
2250 Err(_) => break,
2251 Ok(bytes) => {
2252 b_written += bytes;
2253 if bytes < BLOCK_SIZE.try_into().unwrap() {
2254 break;
2255 }
2256 }
2257 };
2258 }
2259 assert_eq!(
2261 b.file_proxy
2262 .write(&zeros)
2263 .await
2264 .expect("Failed write message")
2265 .expect_err("Write should have been limited"),
2266 Status::NO_SPACE.into_raw()
2267 );
2268
2269 assert!(BLOCK_SIZE as u64 * BLOCK_COUNT as u64 - BYTES_LIMIT >= b_written);
2271 assert!(BLOCK_SIZE as u64 * BLOCK_COUNT as u64 - BYTES_LIMIT <= a_written);
2273
2274 a.file_proxy.close().await.unwrap().expect("Failed to close file");
2275 b.file_proxy.close().await.unwrap().expect("Failed to close file");
2276 volumes_directory.terminate().await;
2277 std::mem::drop(volumes_directory);
2278 filesystem.close().await.expect("close filesystem failed");
2279 }
2280
2281 #[fuchsia::test(threads = 10)]
2282 async fn test_profile_start() {
2283 const PREMOUNT_BLOB: &str = "premount_blob";
2284 const PREMOUNT_NOBLOB: &str = "premount_noblob";
2285 const LIVE_BLOB: &str = "live_blob";
2286 const LIVE_NOBLOB: &str = "live_noblob";
2287
2288 const RECORDING_NAME: &str = "foo";
2289
2290 let device = {
2291 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2292 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2293 let blob_resupplied_count =
2294 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2295 let volumes_directory = VolumesDirectory::new(
2296 root_volume(filesystem.clone()).await.unwrap(),
2297 Weak::new(),
2298 None,
2299 blob_resupplied_count,
2300 MemoryPressureConfig::default(),
2301 )
2302 .await
2303 .unwrap();
2304 volumes_directory
2305 .create_and_mount_volume(PREMOUNT_BLOB, None, true, None)
2306 .await
2307 .unwrap();
2308 volumes_directory
2309 .create_and_mount_volume(PREMOUNT_NOBLOB, None, false, None)
2310 .await
2311 .unwrap();
2312 volumes_directory.create_and_mount_volume(LIVE_BLOB, None, true, None).await.unwrap();
2313 volumes_directory
2314 .create_and_mount_volume(LIVE_NOBLOB, None, false, None)
2315 .await
2316 .unwrap();
2317
2318 volumes_directory.terminate().await;
2319 std::mem::drop(volumes_directory);
2320 filesystem.close().await.expect("Filesystem close");
2321 filesystem.take_device().await
2322 };
2323
2324 device.ensure_unique();
2325 device.reopen(false);
2326 let device = {
2327 let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
2328 let blob_resupplied_count =
2329 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2330 let volumes_directory = VolumesDirectory::new(
2331 root_volume(filesystem.clone()).await.unwrap(),
2332 Weak::new(),
2333 None,
2334 blob_resupplied_count,
2335 MemoryPressureConfig::default(),
2336 )
2337 .await
2338 .unwrap();
2339
2340 let _premount_blob = volumes_directory
2342 .mount_volume(PREMOUNT_BLOB, None, true)
2343 .await
2344 .expect("Reopen volume");
2345 let _premount_noblob = volumes_directory
2346 .mount_volume(PREMOUNT_NOBLOB, None, false)
2347 .await
2348 .expect("Reopen volume");
2349
2350 volumes_directory
2353 .clone()
2354 .record_and_replay_profile(None, RECORDING_NAME.to_owned(), 600)
2355 .await
2356 .expect("Recording");
2357
2358 let _live_blob =
2360 volumes_directory.mount_volume(LIVE_BLOB, None, true).await.expect("Reopen volume");
2361 let _live_noblob = volumes_directory
2362 .mount_volume(LIVE_NOBLOB, None, false)
2363 .await
2364 .expect("Reopen volume");
2365
2366 volumes_directory.stop_profile_tasks().await;
2368
2369 volumes_directory.terminate().await;
2370 std::mem::drop(volumes_directory);
2371 filesystem.close().await.expect("Filesystem close");
2372 filesystem.take_device().await
2373 };
2374
2375 device.ensure_unique();
2376 device.reopen(false);
2377 let filesystem = FxFilesystem::open(device as DeviceHolder).await.unwrap();
2378 {
2379 let blob_resupplied_count =
2380 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2381 let volumes_directory = VolumesDirectory::new(
2382 root_volume(filesystem.clone()).await.unwrap(),
2383 Weak::new(),
2384 None,
2385 blob_resupplied_count,
2386 MemoryPressureConfig::default(),
2387 )
2388 .await
2389 .unwrap();
2390
2391 let _premount_blob = volumes_directory
2392 .mount_volume(PREMOUNT_BLOB, None, true)
2393 .await
2394 .expect("Reopen volume");
2395 let _premount_noblob = volumes_directory
2396 .mount_volume(PREMOUNT_NOBLOB, None, false)
2397 .await
2398 .expect("Reopen volume");
2399 let _live_blob =
2400 volumes_directory.mount_volume(LIVE_BLOB, None, true).await.expect("Reopen volume");
2401 let _live_noblob = volumes_directory
2402 .mount_volume(LIVE_NOBLOB, None, false)
2403 .await
2404 .expect("Reopen volume");
2405
2406 volumes_directory
2408 .delete_profile(PREMOUNT_BLOB, RECORDING_NAME)
2409 .await
2410 .expect("Finding profile to delete.");
2411 volumes_directory
2412 .delete_profile(PREMOUNT_NOBLOB, RECORDING_NAME)
2413 .await
2414 .expect("Finding profile to delete.");
2415 volumes_directory
2416 .delete_profile(LIVE_BLOB, RECORDING_NAME)
2417 .await
2418 .expect("Finding profile to delete.");
2419 volumes_directory
2420 .delete_profile(LIVE_NOBLOB, RECORDING_NAME)
2421 .await
2422 .expect("Finding profile to delete.");
2423
2424 volumes_directory.terminate().await;
2425 }
2426
2427 filesystem.close().await.expect("Filesystem close");
2428 }
2429
2430 #[fuchsia::test(threads = 10)]
2431 async fn test_profile_stop() {
2432 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2433 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2434 let blob_resupplied_count =
2435 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2436 let volumes_directory = VolumesDirectory::new(
2437 root_volume(filesystem.clone()).await.unwrap(),
2438 Weak::new(),
2439 None,
2440 blob_resupplied_count,
2441 MemoryPressureConfig::default(),
2442 )
2443 .await
2444 .unwrap();
2445 let volume =
2446 volumes_directory.create_and_mount_volume("foo", None, true, None).await.unwrap();
2447
2448 volumes_directory
2450 .clone()
2451 .record_and_replay_profile(None, "foo".to_owned(), 0)
2452 .await
2453 .expect("Recording");
2454
2455 while volumes_directory.delete_profile("foo", "foo").await.is_err() {
2457 fasync::Timer::new(Duration::from_millis(10)).await;
2458 }
2459
2460 std::mem::drop(volume);
2461 volumes_directory.terminate().await;
2462 std::mem::drop(volumes_directory);
2463 filesystem.close().await.expect("Filesystem close");
2464 }
2465
2466 #[fuchsia::test(threads = 10)]
2467 async fn test_delete_profile() {
2468 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2469 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2470 let blob_resupplied_count =
2471 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2472 let volumes_directory = VolumesDirectory::new(
2473 root_volume(filesystem.clone()).await.unwrap(),
2474 Weak::new(),
2475 None,
2476 blob_resupplied_count,
2477 MemoryPressureConfig::default(),
2478 )
2479 .await
2480 .unwrap();
2481 let volume =
2482 volumes_directory.create_and_mount_volume("foo", None, true, None).await.unwrap();
2483
2484 volumes_directory
2485 .clone()
2486 .record_and_replay_profile(None, "foo".to_owned(), 600)
2487 .await
2488 .expect("Recording");
2489
2490 assert_eq!(
2492 volumes_directory.delete_profile("foo", "foo").await.expect_err("File shouldn't exist"),
2493 Status::SHOULD_WAIT
2494 );
2495
2496 volumes_directory.stop_profile_tasks().await;
2497
2498 assert_eq!(
2500 volumes_directory.delete_profile("bar", "foo").await.expect_err("File shouldn't exist"),
2501 Status::NOT_FOUND
2502 );
2503
2504 assert_eq!(
2506 volumes_directory.delete_profile("foo", "bar").await.expect_err("File shouldn't exist"),
2507 Status::NOT_FOUND
2508 );
2509
2510 volumes_directory.delete_profile("foo", "foo").await.expect("Deleting");
2512
2513 assert_eq!(
2515 volumes_directory.delete_profile("foo", "foo").await.expect_err("File shouldn't exist"),
2516 Status::NOT_FOUND
2517 );
2518
2519 std::mem::drop(volume);
2520 volumes_directory.terminate().await;
2521 std::mem::drop(volumes_directory);
2522 filesystem.close().await.expect("Filesystem close");
2523 }
2524
2525 #[fuchsia::test(threads = 10)]
2526 async fn test_profile_start_single_volume() {
2527 const TEST_VOLUME: &str = "test_1234";
2528 const TEST_RECORDING: &str = "test_5678";
2529 let crypt = Arc::new(new_insecure_crypt()) as Arc<dyn Crypt>;
2530
2531 let fixture = TestFixture::new().await;
2532 {
2533 let volumes_directory = fixture.volumes_directory();
2534 assert_eq!(
2536 volumes_directory
2537 .record_and_replay_profile(
2538 Some(TEST_VOLUME.to_owned()),
2539 TEST_RECORDING.to_owned(),
2540 1
2541 )
2542 .await
2543 .expect_err("Volumes doesn't exist yet"),
2544 Status::NOT_FOUND
2545 );
2546
2547 {
2549 let volume = volumes_directory
2550 .create_and_mount_volume(TEST_VOLUME, Some(crypt.clone()), false, None)
2551 .await
2552 .unwrap();
2553 volumes_directory
2554 .lock()
2555 .await
2556 .unmount(volume.volume().store().store_object_id())
2557 .await
2558 .expect("unmount failed");
2559 }
2560
2561 assert_eq!(
2563 volumes_directory
2564 .record_and_replay_profile(
2565 Some(TEST_VOLUME.to_owned()),
2566 TEST_RECORDING.to_owned(),
2567 1
2568 )
2569 .await
2570 .expect_err("Volumes doesn't exist yet"),
2571 Status::UNAVAILABLE
2572 );
2573
2574 let volume = volumes_directory
2576 .mount_volume(TEST_VOLUME, Some(crypt.clone()), false)
2577 .await
2578 .expect("Remount volume");
2579 volumes_directory
2580 .record_and_replay_profile(
2581 Some(TEST_VOLUME.to_owned()),
2582 TEST_RECORDING.to_owned(),
2583 1,
2584 )
2585 .await
2586 .expect("Starting recording");
2587
2588 volumes_directory.stop_profile_tasks().await;
2590 {
2591 let profile_dir = volume.volume().get_profile_directory().await.unwrap();
2592 assert!(profile_dir.lookup(TEST_RECORDING).await.unwrap().is_some());
2593 }
2594
2595 {
2597 let profile_dir = fixture.volume().volume().get_profile_directory().await.unwrap();
2598 assert!(profile_dir.lookup(TEST_RECORDING).await.unwrap().is_none());
2599 }
2600 }
2601 fixture.close().await;
2602 }
2603
2604 #[fuchsia::test(threads = 10)]
2605 async fn test_delete_volume_while_flushing() {
2606 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2607 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2608 let blob_resupplied_count =
2609 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2610 let volumes_directory = VolumesDirectory::new(
2611 root_volume(filesystem.clone()).await.unwrap(),
2612 Weak::new(),
2613 None,
2614 blob_resupplied_count,
2615 MemoryPressureConfig::default(),
2616 )
2617 .await
2618 .unwrap();
2619 let name = "vol";
2620 let volume =
2621 volumes_directory.create_and_mount_volume(name, None, false, None).await.unwrap();
2622 let mut transaction = filesystem
2623 .root_store()
2624 .new_transaction(
2625 lock_keys![LockKey::object(
2626 volume.volume().store().store_object_id(),
2627 volume.root_dir().directory().object_id()
2628 )],
2629 Options::default(),
2630 )
2631 .await
2632 .unwrap();
2633 volume
2634 .root_dir()
2635 .directory()
2636 .create_child_file(&mut transaction, "foo")
2637 .await
2638 .expect("create_child_file failed");
2639 transaction.commit().await.expect("commit failed");
2640 volumes_directory
2641 .lock()
2642 .await
2643 .unmount(volume.volume().store().store_object_id())
2644 .await
2645 .expect("unmount failed");
2646
2647 let filesystem_clone = filesystem.clone();
2648 let filesystem_clone2 = filesystem.clone();
2649 let volumes_directory_clone1 = volumes_directory.clone();
2650 let volumes_directory_clone2 = volumes_directory.clone();
2651 let root_store_object_id = filesystem.root_store().store_object_id();
2652 let store_info_object_id = volume.volume().store().store_info_handle_object_id().unwrap();
2653 join!(
2654 async move {
2655 let _guard = filesystem_clone2
2660 .lock_manager()
2661 .read_lock(lock_keys![LockKey::object(
2662 root_store_object_id,
2663 store_info_object_id,
2664 )])
2665 .await;
2666 fasync::Timer::new(Duration::from_millis(200)).await;
2667 },
2668 async move {
2669 filesystem_clone.journal().force_compact().await.expect("Compact failed");
2670 },
2671 async move {
2672 if let Err(e) = volumes_directory_clone1.remove_volume(name).await {
2673 if !FxfsError::NotFound.matches(&e) {
2674 panic!("remove_volume failed: {e:?}");
2675 }
2676 }
2677 },
2678 async move {
2679 if let Err(e) = volumes_directory_clone2.remove_volume(name).await {
2680 if !FxfsError::NotFound.matches(&e) {
2681 panic!("remove_volume failed: {e:?}");
2682 }
2683 }
2684 },
2685 );
2686 volumes_directory.terminate().await;
2687 std::mem::drop(volumes_directory);
2688 filesystem.close().await.expect("Filesystem close");
2689 }
2690
2691 #[fuchsia::test(threads = 10)]
2694 async fn test_flush_before_mark_dirty_under_critical_memory_pressure() {
2695 let fixture = TestFixture::new().await;
2696 let _ = fixture
2698 .memory_pressure_proxy()
2699 .on_level_changed(MemoryPressureLevel::Critical)
2700 .await
2701 .expect("memory pressure FIDL");
2702 fixture.volumes_directory().max_dirty_bytes_when_critical.store(1, Ordering::Relaxed);
2703
2704 let root = fixture.root();
2705 let file = open_file_checked(
2706 &root,
2707 "foo",
2708 fio::Flags::FLAG_MAYBE_CREATE
2709 | fio::PERM_READABLE
2710 | fio::PERM_WRITABLE
2711 | fio::Flags::PROTOCOL_FILE,
2712 &Default::default(),
2713 )
2714 .await;
2715
2716 file.resize((zx::system_get_page_size() * 2).into())
2717 .await
2718 .expect("resize (FIDL)")
2719 .expect("resize failed");
2720 file.sync().await.expect("Failed to make sync call").expect("sync failed");
2738
2739 let vmo = file
2740 .get_backing_memory(fio::VmoFlags::READ | fio::VmoFlags::WRITE)
2741 .await
2742 .expect("get_backing_memory (FIDL)")
2743 .expect("get_backing_memory");
2744
2745 let buf = [0xAAu8];
2746 vmo.write(&buf, 0).expect("Writing to create dirty bytes");
2748 let before = fixture.volumes_directory().pager_dirty_bytes_count.load();
2749 vmo.write(&buf, zx::system_get_page_size().into())
2750 .expect("Writing to force a flush during mark_dirty");
2751 assert_eq!(fixture.volumes_directory().pager_dirty_bytes_count.load(), before,);
2754
2755 fixture.close().await;
2756 }
2757
2758 #[fuchsia::test(threads = 10)]
2759 async fn test_delete_crypt_for_volume() {
2760 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2761 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2762 let store_id;
2763 {
2764 let blob_resupplied_count =
2765 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2766 let volumes_directory = VolumesDirectory::new(
2767 root_volume(filesystem.clone()).await.unwrap(),
2768 Weak::new(),
2769 None,
2770 blob_resupplied_count,
2771 MemoryPressureConfig::default(),
2772 )
2773 .await
2774 .unwrap();
2775 let name = "vol";
2776 let crypt = Arc::new(new_insecure_crypt());
2777 let volume = volumes_directory
2778 .create_and_mount_volume(name, Some(crypt.clone()), false, None)
2779 .await
2780 .unwrap();
2781 store_id = volume.volume().store().store_object_id();
2782 let mut transaction = filesystem
2784 .root_store()
2785 .new_transaction(
2786 lock_keys![LockKey::object(
2787 volume.volume().store().store_object_id(),
2788 volume.root_dir().directory().object_id()
2789 )],
2790 Options::default(),
2791 )
2792 .await
2793 .unwrap();
2794 volume
2795 .root_dir()
2796 .directory()
2797 .create_child_file(&mut transaction, "foo")
2798 .await
2799 .expect("create_child_file failed");
2800 transaction.commit().await.expect("commit failed");
2801
2802 let (volume_dir_proxy, dir_server_end) =
2803 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2804 volumes_directory
2805 .serve_volume(&volume, dir_server_end, false)
2806 .expect("serve_volume failed");
2807 let (root_dir, root_server_end) =
2808 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2809 volume_dir_proxy
2810 .open(
2811 "root",
2812 fio::PERM_READABLE | fio::PERM_WRITABLE | fio::Flags::PROTOCOL_DIRECTORY,
2813 &Default::default(),
2814 root_server_end.into_channel(),
2815 )
2816 .expect("Failed to open volume root");
2817
2818 let filesystem_clone = filesystem.clone();
2819 join!(
2820 async move {
2821 filesystem_clone.journal().force_compact().await.expect("Compact failed");
2822 },
2823 async move {
2824 let mut i = 0;
2825 while let Ok(_) = fuchsia_fs::directory::open_file(
2826 &root_dir,
2827 &format!("foo{i}"),
2828 fio::Flags::FLAG_MAYBE_CREATE | fio::PERM_READABLE,
2829 )
2830 .await
2831 {
2832 i += 1;
2833 }
2834 },
2835 async move {
2836 crypt.shutdown();
2837 },
2838 );
2839
2840 let (admin_proxy, server_end) =
2842 fidl::endpoints::create_proxy::<fidl_fuchsia_fs::AdminMarker>();
2843 volume_dir_proxy
2844 .open(
2845 &format!("svc/{}", fidl_fuchsia_fs::AdminMarker::PROTOCOL_NAME),
2846 fio::Flags::PROTOCOL_SERVICE,
2847 &Default::default(),
2848 server_end.into(),
2849 )
2850 .expect("Failed to open Admin connection");
2851 admin_proxy.shutdown().await.expect("shutdown failed");
2852
2853 volumes_directory.terminate().await;
2854 }
2855 filesystem.close().await.expect("Filesystem close");
2856 let device = filesystem.take_device().await;
2857 device.reopen(false);
2858 let filesystem = FxFilesystem::open(device).await.expect("open failed");
2859 let options = FsckOptions { fail_on_warning: true, ..Default::default() };
2860 fsck_with_options(filesystem.clone(), &options).await.expect("fsck failed");
2861 fsck_volume_with_options(
2862 filesystem.as_ref(),
2863 &options,
2864 store_id,
2865 Some(Arc::new(new_insecure_crypt())),
2866 )
2867 .await
2868 .expect("fsck_volume failed");
2869 filesystem.close().await.expect("Filesystem close");
2870 }
2871
2872 #[fuchsia::test(threads = 10)]
2874 async fn test_volume_installation() {
2875 let fixture = TestFixture::open(
2876 DeviceHolder::new(FakeDevice::new(1024, 4096)),
2877 TestFixtureOptions { format: true, encrypted: false, ..Default::default() },
2878 )
2879 .await;
2880
2881 {
2883 let file = open_file_checked(
2884 fixture.root(),
2885 "foo",
2886 fio::Flags::PROTOCOL_FILE | fio::Flags::FLAG_MUST_CREATE | fio::PERM_WRITABLE,
2887 &Default::default(),
2888 )
2889 .await;
2890 file.write("Hello, world!".as_bytes()).await.unwrap().expect("write failed");
2891 };
2892
2893 let image = {
2895 let inner_fixture = TestFixture::open(
2896 DeviceHolder::new(FakeDevice::new(512, 4096)),
2897 TestFixtureOptions { format: true, encrypted: false, ..Default::default() },
2898 )
2899 .await;
2900 let file = open_file_checked(
2901 inner_fixture.root(),
2902 "bar",
2903 fio::Flags::PROTOCOL_FILE | fio::Flags::FLAG_MUST_CREATE | fio::PERM_WRITABLE,
2904 &Default::default(),
2905 )
2906 .await;
2907 file.write("Well, this is new...".as_bytes()).await.unwrap().expect("write failed");
2908 file.close().await.unwrap().expect("close error");
2909 inner_fixture.close().await
2910 };
2911
2912 {
2914 let (src_out_dir, server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
2915 fixture
2916 .volumes_directory()
2917 .create_and_serve_volume("src", server_end, Default::default(), Default::default())
2918 .await
2919 .unwrap();
2920 let src_root = open_dir_checked(
2921 &src_out_dir,
2922 "root",
2923 fio::PERM_READABLE | fio::PERM_WRITABLE,
2924 Default::default(),
2925 )
2926 .await;
2927 let file = open_file_checked(
2928 &src_root,
2929 "image",
2930 fio::Flags::PROTOCOL_FILE | fio::Flags::FLAG_MUST_CREATE | fio::PERM_WRITABLE,
2931 &Default::default(),
2932 )
2933 .await;
2934 write_image_to_file(image, file).await;
2935 };
2936
2937 assert!(
2939 fixture.volumes_directory().install_volume("src", "image", "vol").await.is_err(),
2940 "volume installation should fail while either src/dst is still mounted"
2941 );
2942
2943 let device = fixture.close().await;
2946 let fs = FxFilesystem::open(device).await.unwrap();
2947 {
2948 let root = root_volume(fs.clone()).await.unwrap();
2949 root.install_volume("src", "image", "vol").await.unwrap();
2950 }
2951 fs.close().await.unwrap();
2952 let device = fs.take_device().await;
2953 device.reopen(true);
2954 let fixture = TestFixture::open(
2955 device,
2956 TestFixtureOptions { encrypted: false, format: false, ..Default::default() },
2957 )
2958 .await;
2959
2960 assert!(
2962 fixture.volumes_directory().mount_volume("src", None, false).await.is_err(),
2963 "src volume should be deleted after installation"
2964 );
2965 assert!(
2966 testing::open_file(
2967 fixture.volume_out_dir(),
2968 "foo",
2969 fio::PERM_READABLE,
2970 &Default::default()
2971 )
2972 .await
2973 .is_err(),
2974 "foo should be deleted after installation"
2975 );
2976
2977 let file =
2979 open_file_checked(fixture.root(), "bar", fio::PERM_READABLE, &Default::default()).await;
2980 let data = file.read(fio::MAX_TRANSFER_SIZE).await.unwrap().expect("read failed");
2981 assert_eq!(String::from_utf8(data).unwrap(), "Well, this is new...");
2982 file.close().await.unwrap().unwrap();
2983
2984 fixture.close().await;
2985 }
2986
2987 #[fuchsia::test]
2988 async fn test_create_with_low_32_bit_ids() {
2989 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
2990 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
2991 let blob_resupplied_count =
2992 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
2993
2994 {
2995 let volumes_directory = VolumesDirectory::new(
2996 root_volume(filesystem.clone()).await.unwrap(),
2997 Weak::new(),
2998 None,
2999 blob_resupplied_count,
3000 MemoryPressureConfig::default(),
3001 )
3002 .await
3003 .unwrap();
3004
3005 let mut guard = volumes_directory.lock().await;
3006
3007 let vol = guard
3008 .create_or_mount_volume(
3009 "low_32",
3010 None,
3011 Mode::Create { guid: None, low_32_bit_object_ids: true },
3012 false,
3013 )
3014 .await
3015 .expect("create volume failed");
3016
3017 let root_dir = vol.volume().store().root_directory_object_id();
3018 let root_dir = fxfs::object_store::Directory::open(vol.volume().store(), root_dir)
3019 .await
3020 .expect("open failed");
3021
3022 let mut transaction = filesystem
3023 .root_store()
3024 .new_transaction(
3025 lock_keys![LockKey::object(
3026 vol.volume().store().store_object_id(),
3027 root_dir.object_id()
3028 )],
3029 Options::default(),
3030 )
3031 .await
3032 .expect("new_transaction failed");
3033
3034 let object = root_dir
3035 .create_child_file(&mut transaction, "test")
3036 .await
3037 .expect("create_child_file failed");
3038
3039 assert!(object.object_id() < 1 << 32);
3041 transaction.commit().await.expect("commit failed");
3042 };
3043
3044 filesystem.close().await.expect("close filesystem failed");
3045
3046 let device = filesystem.take_device().await;
3048 device.reopen(false);
3049 let filesystem = FxFilesystem::open(device).await.unwrap();
3050 let blob_resupplied_count =
3051 Arc::new(PageRefaultCounter::new().expect("Failed to create PageRefaultCounter"));
3052 let volumes_directory = VolumesDirectory::new(
3053 root_volume(filesystem.clone()).await.unwrap(),
3054 Weak::new(),
3055 None,
3056 blob_resupplied_count,
3057 MemoryPressureConfig::default(),
3058 )
3059 .await
3060 .unwrap();
3061
3062 {
3064 let mut guard = volumes_directory.lock().await;
3065
3066 let vol = guard
3067 .create_or_mount_volume("low_32", None, Mode::Mount, false)
3068 .await
3069 .expect("mount volume failed");
3070
3071 let root_dir = vol.volume().store().root_directory_object_id();
3072 let root_dir = fxfs::object_store::Directory::open(vol.volume().store(), root_dir)
3073 .await
3074 .expect("open failed");
3075
3076 let mut transaction = filesystem
3077 .root_store()
3078 .new_transaction(
3079 lock_keys![LockKey::object(
3080 vol.volume().store().store_object_id(),
3081 root_dir.object_id()
3082 )],
3083 Options::default(),
3084 )
3085 .await
3086 .expect("new_transaction failed");
3087
3088 let object = root_dir
3089 .create_child_file(&mut transaction, "test2")
3090 .await
3091 .expect("create_child_file failed");
3092
3093 assert!(object.object_id() < 1 << 32);
3094 transaction.commit().await.expect("commit failed");
3095 }
3096
3097 filesystem.close().await.expect("close filesystem failed");
3098 }
3099
3100 #[fuchsia::test(threads = 10)]
3101 async fn test_race_unmount_and_flush_with_crypt_error() {
3102 let device = DeviceHolder::new(FakeDevice::new(8192, 512));
3103 let filesystem = FxFilesystem::new_empty(device).await.unwrap();
3104 let blob_resupplied_count = Arc::new(PageRefaultCounter::new().unwrap());
3105 let volumes_directory = VolumesDirectory::new(
3106 root_volume(filesystem.clone()).await.unwrap(),
3107 Weak::new(),
3108 None,
3109 blob_resupplied_count,
3110 MemoryPressureConfig::default(),
3111 )
3112 .await
3113 .unwrap();
3114
3115 let crypt_service = Arc::new(fxfs_crypt::CryptService::new());
3116 crypt_service
3117 .add_wrapping_key(0, fxfs_insecure_crypto::DATA_KEY.to_vec())
3118 .expect("add_wrapping_key failed");
3119 crypt_service
3120 .add_wrapping_key(1, fxfs_insecure_crypto::METADATA_KEY.to_vec())
3121 .expect("add_wrapping_key failed");
3122 crypt_service.set_active_key(KeyPurpose::Data, 0).expect("set_active_key failed");
3123 crypt_service.set_active_key(KeyPurpose::Metadata, 1).expect("set_active_key failed");
3124
3125 for _ in 0..20 {
3126 let (client, mut stream) = create_request_stream::<fidl_fuchsia_fxfs::CryptMarker>();
3127 let (close_tx, mut close_rx) = futures::channel::oneshot::channel::<()>();
3128
3129 let crypt_task = fasync::Task::spawn(async move {
3130 loop {
3131 futures::select! {
3132 _ = close_rx => return, request = stream.try_next() => {
3134 match request {
3135 Ok(Some(CryptRequest::CreateKey { responder, .. })) => {
3136 responder.send(Ok((&[0; 16], &[0; 48], &[0; 32]))).unwrap();
3137 }
3138 Ok(Some(CryptRequest::CreateKeyWithId {
3139 wrapping_key_id, responder, .. })) => {
3140 let key = WrappedKey::Fxfs(FxfsKey {
3141 wrapping_key_id,
3142 wrapped_key: [0u8; 48],
3143 });
3144 responder.send(Ok((&key, &[0; 32]))).unwrap();
3145 }
3146 Ok(Some(CryptRequest::UnwrapKey { responder, .. })) => {
3147 responder.send(Ok(&vec![0; 32])).unwrap();
3148 }
3149 _ => return,
3150 }
3151 }
3152 }
3153 }
3154 });
3155
3156 let volume = volumes_directory
3157 .create_and_mount_volume(
3158 "encrypted",
3159 Some(Arc::new(RemoteCrypt::new(client))),
3160 false,
3161 None,
3162 )
3163 .await
3164 .unwrap();
3165
3166 {
3168 let mut transaction = filesystem
3169 .root_store()
3170 .new_transaction(
3171 lock_keys![LockKey::object(
3172 volume.volume().store().store_object_id(),
3173 volume.root_dir().directory().object_id()
3174 )],
3175 Options::default(),
3176 )
3177 .await
3178 .unwrap();
3179 volume
3180 .root_dir()
3181 .directory()
3182 .create_child_file(&mut transaction, "foo")
3183 .await
3184 .expect("create_child_file failed");
3185 transaction.commit().await.expect("commit failed");
3186 }
3187
3188 let (dir_proxy, dir_server_end) =
3189 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
3190 volumes_directory.serve_volume(&volume, dir_server_end, false).unwrap();
3191 volumes_directory.lock().await.auto_unmount(volume.volume().store().store_object_id());
3192
3193 let _ = close_tx.send(());
3195
3196 let filesystem_clone = filesystem.clone();
3197 let compact_task = fasync::Task::spawn(async move {
3198 filesystem_clone.object_manager().flush().await.expect("flush failed");
3202 });
3203
3204 std::mem::drop(dir_proxy);
3206
3207 let store_id = volume.volume().store().store_object_id();
3209 loop {
3210 {
3211 let guard = volumes_directory.lock().await;
3212 if !guard.mounted_volumes.contains_key(&store_id) {
3213 break;
3214 }
3215 }
3216 fasync::Timer::new(Duration::from_millis(10)).await;
3217 }
3218
3219 join!(compact_task, crypt_task);
3220 volumes_directory.remove_volume("encrypted").await.expect("remove_volume failed");
3221 }
3222 volumes_directory.terminate().await;
3223 filesystem.close().await.expect("close filesystem failed");
3224 }
3225}