1use crate::fuchsia::debug::{create_debug_directory, handle_debug_request};
6use crate::fuchsia::errors::map_to_status;
7use crate::fuchsia::memory_pressure::MemoryPressureMonitor;
8use crate::fuchsia::power::FuchsiaPowerManager;
9use crate::fuchsia::volume::MemoryPressureConfig;
10use crate::fuchsia::volumes_directory::VolumesDirectory;
11use anyhow::{Context, Error, bail};
12use async_trait::async_trait;
13use block_client::{BlockClient as _, RemoteBlockClient};
14use fidl::endpoints::{ClientEnd, DiscoverableProtocolMarker, RequestStream};
15use fidl_fuchsia_fs::{AdminMarker, AdminRequest, AdminRequestStream};
16use fidl_fuchsia_fs_startup::{
17 CheckOptions, StartOptions, StartupMarker, StartupRequest, StartupRequestStream, VolumesMarker,
18 VolumesRequest, VolumesRequestStream,
19};
20use fidl_fuchsia_fxfs::{
21 DebugMarker, DebugRequestStream, VolumeInstallerMarker, VolumeInstallerRequest,
22 VolumeInstallerRequestStream,
23};
24use fidl_fuchsia_io as fio;
25use fidl_fuchsia_memory_attribution as fattribution;
26use fidl_fuchsia_process_lifecycle::{LifecycleRequest, LifecycleRequestStream};
27use fidl_fuchsia_storage_block::BlockMarker;
28use fs_inspect::{FsInspect, FsInspectTree, InfoData, UsageData};
29use fuchsia_async as fasync;
30use fuchsia_component_client::connect_to_protocol;
31use futures::TryStreamExt;
32use futures::lock::Mutex;
33use fxfs::filesystem::{FxFilesystem, FxFilesystemBuilder, MIN_BLOCK_SIZE, OpenFxFilesystem, mkfs};
34use fxfs::log::*;
35use fxfs::object_store::volume::root_volume;
36use fxfs::serialized_types::LATEST_VERSION;
37use fxfs::{fsck, metrics};
38use fxfs_trace::{TraceFutureExt, trace_future_args};
39use refaults_vmo::PageRefaultCounter;
40use std::ops::Deref;
41use std::sync::{Arc, Weak};
42use storage_device::DeviceHolder;
43use storage_device::block_device::BlockDevice;
44use vfs::directory::helper::DirectlyMutable;
45use vfs::execution_scope::ExecutionScope;
46
47const FXFS_INFO_NAME: &'static str = "fxfs";
48
49pub fn map_to_raw_status(e: Error) -> zx::sys::zx_status_t {
50 map_to_status(e).into_raw()
51}
52
53pub async fn new_block_client(remote: ClientEnd<BlockMarker>) -> Result<RemoteBlockClient, Error> {
54 Ok(RemoteBlockClient::new(remote.into_proxy()).await?)
55}
56
57pub struct Component {
59 state: futures::lock::Mutex<State>,
60
61 scope: ExecutionScope,
63
64 export_dir: Arc<vfs::directory::immutable::Simple>,
66}
67
68struct InspectedFxFilesystem(OpenFxFilesystem, u64);
70
71impl From<OpenFxFilesystem> for InspectedFxFilesystem {
72 fn from(fs: OpenFxFilesystem) -> Self {
73 Self(fs, zx::Event::create().koid().unwrap().raw_koid())
74 }
75}
76
77impl Deref for InspectedFxFilesystem {
78 type Target = Arc<FxFilesystem>;
79 fn deref(&self) -> &Self::Target {
80 &self.0.deref()
81 }
82}
83
84#[async_trait]
85impl FsInspect for InspectedFxFilesystem {
86 fn get_info_data(&self) -> InfoData {
87 let earliest_version = self.0.super_block_header().earliest_version;
88 InfoData {
89 id: self.1,
90 fs_type: fidl_fuchsia_fs::VfsType::Fxfs.into_primitive().into(),
91 name: FXFS_INFO_NAME.into(),
92 version_major: LATEST_VERSION.major.into(),
93 version_minor: LATEST_VERSION.minor.into(),
94 block_size: self.0.block_size() as u64,
95 max_filename_length: fio::MAX_NAME_LENGTH,
96 oldest_version: Some(format!("{}.{}", earliest_version.major, earliest_version.minor)),
97 }
98 }
99
100 async fn get_usage_data(&self) -> UsageData {
101 let info = self.0.get_info();
102 UsageData {
103 total_bytes: info.total_bytes,
104 used_bytes: info.used_bytes,
105 total_nodes: 0,
107 used_nodes: 0,
108 }
109 }
110}
111
112struct RunningState {
113 fs: Arc<InspectedFxFilesystem>,
116 volumes: Arc<VolumesDirectory>,
117 _inspect_tree: Arc<FsInspectTree>,
118}
119
120enum State {
121 ComponentStarted,
122 Running(RunningState),
123}
124
125impl State {
126 async fn stop(&mut self, outgoing_dir: &vfs::directory::immutable::Simple) {
127 if let State::Running(RunningState { fs, volumes, .. }) =
128 std::mem::replace(self, State::ComponentStarted)
129 {
130 info!("Stopping Fxfs runtime; remaining connections will be forcibly closed");
131 let _ = outgoing_dir.remove_entry("volumes", false);
132 let _ = outgoing_dir.remove_entry("debug", false);
133 volumes.terminate().await;
134 let _ = fs.deref().close().await;
135 }
136 }
137}
138
139impl Component {
140 pub fn new() -> Arc<Self> {
141 Arc::new(Self {
142 state: Mutex::new(State::ComponentStarted),
143 scope: ExecutionScope::new(),
144 export_dir: vfs::directory::immutable::simple(),
145 })
146 }
147
148 pub async fn run(
150 self: Arc<Self>,
151 outgoing_dir: zx::Channel,
152 lifecycle_channel: Option<zx::Channel>,
153 ) -> Result<(), Error> {
154 let svc_dir = vfs::directory::immutable::simple();
155 self.export_dir.add_entry("svc", svc_dir.clone()).expect("Unable to create svc dir");
156
157 let weak = Arc::downgrade(&self);
158 svc_dir.add_entry(
159 StartupMarker::PROTOCOL_NAME,
160 vfs::service::host(move |requests| {
161 let weak = weak.clone();
162 async move {
163 if let Some(me) = weak.upgrade() {
164 let _ = me.handle_startup_requests(requests).await;
165 }
166 }
167 }),
168 )?;
169 let weak = Arc::downgrade(&self);
170 svc_dir.add_entry(
171 VolumesMarker::PROTOCOL_NAME,
172 vfs::service::host(move |requests| {
173 let weak = weak.clone();
174 async move {
175 if let Some(me) = weak.upgrade() {
176 me.handle_volumes_requests(requests).await;
177 }
178 }
179 }),
180 )?;
181
182 let weak = Arc::downgrade(&self);
183 svc_dir.add_entry(
184 AdminMarker::PROTOCOL_NAME,
185 vfs::service::host(move |requests| {
186 let weak = weak.clone();
187 async move {
188 if let Some(me) = weak.upgrade() {
189 let _ = me.handle_admin_requests(requests).await;
190 }
191 }
192 }),
193 )?;
194
195 let weak = Arc::downgrade(&self);
196 svc_dir.add_entry(
197 VolumeInstallerMarker::PROTOCOL_NAME,
198 vfs::service::host(move |requests| {
199 let weak = weak.clone();
200 async move {
201 if let Some(me) = weak.upgrade() {
202 if let Err(error) = me.handle_volume_installer_requests(requests).await {
203 error!(error:?; "failed to handle VolumeInstaller requests");
204 }
205 }
206 }
207 }),
208 )?;
209
210 let weak = Arc::downgrade(&self);
212 svc_dir.add_entry(
213 DebugMarker::PROTOCOL_NAME,
214 vfs::service::host(move |requests| {
215 let weak = weak.clone();
216 async move {
217 if let Some(me) = weak.upgrade() {
218 let _ = me.handle_debug_requests(requests).await;
219 }
220 }
221 }),
222 )?;
223
224 vfs::directory::serve_on(
225 self.export_dir.clone(),
226 fio::PERM_READABLE | fio::PERM_WRITABLE | fio::PERM_EXECUTABLE,
227 self.scope.clone(),
228 fidl::endpoints::ServerEnd::new(outgoing_dir),
229 );
230
231 if let Some(channel) = lifecycle_channel {
232 let me = self.clone();
233 self.scope.spawn(async move {
234 if let Err(error) = me.handle_lifecycle_requests(channel).await {
235 warn!(error:?; "handle_lifecycle_requests");
236 }
237 });
238 }
239
240 self.scope.wait().await;
241
242 Ok(())
243 }
244
245 async fn handle_startup_requests(&self, mut stream: StartupRequestStream) -> Result<(), Error> {
246 while let Some(request) = stream.try_next().await? {
247 match request {
248 StartupRequest::Start { responder, device, options } => {
249 async move {
250 responder.send(self.handle_start(device, options).await.map_err(|error| {
251 error!(error:?; "handle_start failed");
252 map_to_raw_status(error)
253 }))
254 }
255 .trace(trace_future_args!("Startup::Start"))
256 .await?
257 }
258 StartupRequest::Format { responder, device, .. } => {
259 async move {
260 responder.send(self.handle_format(device).await.map_err(|error| {
261 error!(error:?; "handle_format failed");
262 map_to_raw_status(error)
263 }))
264 }
265 .trace(trace_future_args!("Startup::Format"))
266 .await?
267 }
268 StartupRequest::Check { responder, device, options } => {
269 async move {
270 responder.send(self.handle_check(device, options).await.map_err(|error| {
271 error!(error:?; "handle_check failed");
272 map_to_raw_status(error)
273 }))
274 }
275 .trace(trace_future_args!("Startup::Check"))
276 .await?
277 }
278 }
279 }
280 Ok(())
281 }
282
283 async fn handle_start(
284 &self,
285 device: ClientEnd<BlockMarker>,
286 options: StartOptions,
287 ) -> Result<(), Error> {
288 info!(options:?; "Received start request");
289 let mut state = self.state.lock().await;
290 state.stop(&self.export_dir).await;
294 let client = new_block_client(device).await?;
295
296 assert!(client.block_size() <= zx::system_get_page_size());
298 assert!((zx::system_get_page_size() as u64) == MIN_BLOCK_SIZE);
299
300 let fs = FxFilesystemBuilder::new()
301 .fsck_after_every_transaction(options.fsck_after_every_transaction.unwrap_or(false))
302 .read_only(options.read_only.unwrap_or(false))
303 .inline_crypto_enabled(options.inline_crypto_enabled.unwrap_or(false))
304 .barriers_enabled(options.barriers_enabled.unwrap_or(false))
305 .allow_type3_blobs(options.allow_type3_blobs.unwrap_or(false))
306 .power_manager(FuchsiaPowerManager::new())
307 .open(DeviceHolder::new(
308 BlockDevice::new(client, options.read_only.unwrap_or(false)).await?,
309 ))
310 .await?;
311 let root_volume = root_volume(fs.clone()).await?;
312 let fs: Arc<InspectedFxFilesystem> = Arc::new(fs.into());
313 let weak_fs = Arc::downgrade(&fs) as Weak<dyn FsInspect + Send + Sync>;
314 let inspect_tree =
315 Arc::new(FsInspectTree::new(weak_fs, fuchsia_inspect::component::inspector().root()));
316
317 let blob_resupplied_count = Arc::new(PageRefaultCounter::new()?);
318 connect_to_protocol::<fattribution::PageRefaultSinkMarker>()?
319 .send_page_refault_count(blob_resupplied_count.readonly_vmo()?)?;
320
321 let mem_monitor = match MemoryPressureMonitor::start() {
322 Ok(v) => Some(v),
323 Err(error) => {
324 warn!(
325 error:?;
326 "Failed to connect to memory pressure monitor. Running \
327 without pressure awareness."
328 );
329 None
330 }
331 };
332
333 let volumes = VolumesDirectory::new(
334 root_volume,
335 Arc::downgrade(&inspect_tree),
336 mem_monitor,
337 blob_resupplied_count,
338 MemoryPressureConfig::default(),
339 )
340 .await?;
341
342 self.export_dir.add_entry_may_overwrite(
343 "volumes",
344 volumes.directory_node().clone(),
345 true,
346 )?;
347
348 let debug = create_debug_directory(&**fs, &volumes);
349 self.export_dir.add_entry_may_overwrite("debug", debug, true)?;
350
351 fs.allocator().track_statistics(&metrics::detail(), "allocator");
352 fs.journal().track_statistics(&metrics::detail(), "journal");
353 fs.object_manager().track_statistics(&metrics::detail(), "object_manager");
354
355 let info = fs.get_info();
356 info!(
357 device_size = info.total_bytes,
358 used = info.used_bytes,
359 free = info.total_bytes - info.used_bytes;
360 "Mounted"
361 );
362
363 if let Some(profile_time) = options.startup_profiling_seconds {
364 volumes
366 .record_and_replay_profile(None, ".boot".to_owned(), profile_time)
367 .await
368 .unwrap();
369 }
370
371 *state = State::Running(RunningState { fs, volumes, _inspect_tree: inspect_tree });
372
373 Ok(())
374 }
375
376 async fn handle_format(&self, device: ClientEnd<BlockMarker>) -> Result<(), Error> {
377 let device = DeviceHolder::new(
378 BlockDevice::new(new_block_client(device).await?, false).await?,
379 );
380 mkfs(device).await?;
381 info!("Formatted filesystem");
382 Ok(())
383 }
384
385 async fn handle_check(
386 &self,
387 device: ClientEnd<BlockMarker>,
388 _options: CheckOptions,
389 ) -> Result<(), Error> {
390 let state = self.state.lock().await;
391 let (fs_container, fs) = match *state {
392 State::ComponentStarted => {
393 let client = new_block_client(device).await?;
394 let fs_container = FxFilesystemBuilder::new()
395 .read_only(true)
396 .open(DeviceHolder::new(BlockDevice::new(client, true).await?))
397 .await?;
398 let fs = fs_container.clone();
399 (Some(fs_container), fs)
400 }
401 State::Running(RunningState { ref fs, .. }) => (None, fs.deref().deref().clone()),
402 };
403 let res = fsck::fsck(fs.clone()).await;
404 if let Some(fs_container) = fs_container {
405 let _ = fs_container.close().await;
406 }
407 info!("handle_check for fs: {:?}", res?);
409 Ok(())
410 }
411
412 async fn handle_admin_requests(&self, mut stream: AdminRequestStream) -> Result<(), Error> {
413 while let Some(request) = stream.try_next().await.context("Reading request")? {
414 if self.handle_admin(request).await? {
415 break;
416 }
417 }
418 Ok(())
419 }
420
421 async fn handle_admin(&self, req: AdminRequest) -> Result<bool, Error> {
423 match req {
424 AdminRequest::Shutdown { responder } => {
425 async move {
426 info!("Received shutdown request");
427 self.shutdown().await;
428 responder
429 .send()
430 .unwrap_or_else(|error| warn!(error:?; "Failed to send shutdown response"));
431 }
432 .trace(trace_future_args!("Admin::Shutdown"))
433 .await;
434 return Ok(true);
435 }
436 }
437 }
438
439 async fn handle_debug_requests(&self, mut stream: DebugRequestStream) -> Result<(), Error> {
442 while let Some(request) = stream.try_next().await.context("Reading request")? {
443 let state = self.state.lock().await;
444 let (fs, volumes) = match &*state {
445 State::ComponentStarted => {
446 info!("Debug commands are not valid unless component is started.");
447 bail!("Component not started");
448 }
449 State::Running(RunningState { fs, volumes, .. }) => {
450 (fs.deref().deref().clone(), volumes.clone())
451 }
452 };
453 handle_debug_request(fs, volumes, request).await?;
454 }
455 Ok(())
456 }
457
458 async fn shutdown(&self) {
459 self.state.lock().await.stop(&self.export_dir).await;
460 info!("Filesystem terminated");
461 }
462
463 async fn handle_volumes_requests(&self, mut stream: VolumesRequestStream) {
464 let volumes =
465 if let State::Running(RunningState { volumes, .. }) = &*self.state.lock().await {
466 volumes.clone()
467 } else {
468 let _ = stream.into_inner().0.shutdown_with_epitaph(zx::Status::BAD_STATE);
469 return;
470 };
471 while let Ok(Some(request)) = stream.try_next().await {
472 match request {
473 VolumesRequest::Create {
474 name,
475 outgoing_directory,
476 create_options,
477 mount_options,
478 responder,
479 } => {
480 async {
481 info!(
482 name = name.as_str();
483 "Create {}volume",
484 if mount_options.crypt.is_some() { "encrypted " } else { "" }
485 );
486 responder
487 .send(
488 volumes
489 .create_and_serve_volume(
490 &name,
491 outgoing_directory.into_channel().into(),
492 mount_options,
493 create_options,
494 )
495 .await
496 .map_err(map_to_raw_status),
497 )
498 .unwrap_or_else(
499 |error| warn!(error:?; "Failed to send volume creation response"),
500 );
501 }
502 .trace(trace_future_args!("Volumes::Create"))
503 .await;
504 }
505 VolumesRequest::Remove { name, responder } => {
506 async {
507 info!(name = name.as_str(); "Remove volume");
508 responder
509 .send(volumes.remove_volume(&name).await.map_err(map_to_raw_status))
510 .unwrap_or_else(
511 |error| warn!(error:?; "Failed to send volume removal response"),
512 );
513 }
514 .trace(trace_future_args!("Volumes::Remove"))
515 .await;
516 }
517 VolumesRequest::GetInfo { responder } => {
518 async move {
519 responder.send(Err(zx::Status::NOT_SUPPORTED.into_raw())).unwrap_or_else(
520 |error| warn!(error:?; "Failed to send volume removal response"),
521 )
522 }
523 .trace(trace_future_args!("Volumes::GetInfo"))
524 .await;
525 }
526 }
527 }
528 }
529
530 async fn handle_lifecycle_requests(&self, lifecycle_channel: zx::Channel) -> Result<(), Error> {
531 let mut stream =
532 LifecycleRequestStream::from_channel(fasync::Channel::from_channel(lifecycle_channel));
533 match stream.try_next().await.context("Reading request")? {
534 Some(LifecycleRequest::Stop { .. }) => {
535 info!("Received Lifecycle::Stop request");
536 self.shutdown().await;
537 }
538 None => {}
539 }
540 Ok(())
541 }
542
543 async fn handle_volume_installer_requests(
544 &self,
545 mut stream: VolumeInstallerRequestStream,
546 ) -> Result<(), Error> {
547 let volumes =
548 if let State::Running(RunningState { volumes, .. }) = &*self.state.lock().await {
549 volumes.clone()
550 } else {
551 let _ = stream.into_inner().0.shutdown_with_epitaph(zx::Status::BAD_STATE);
552 return Err(zx::Status::BAD_STATE).context("fxfs component not running");
553 };
554 while let Some(request) = stream.try_next().await.context("reading request")? {
555 match request {
556 VolumeInstallerRequest::Install { src, image_file, dst, responder } => {
557 async {
558 let response = volumes.install_volume(&src, &image_file, &dst).await;
559 responder.send(response.map_err(|error| {
560 error!(error:?; "install failed");
561 map_to_raw_status(error)
562 }))
563 }
564 .trace(trace_future_args!("VolumeInstaller::Install"))
565 .await?;
566 }
567 }
568 }
569 Ok(())
570 }
571}
572
573#[cfg(test)]
574mod tests {
575 use super::{Component, new_block_client};
576 use fidl::endpoints::Proxy;
577 use fidl_fuchsia_fs::AdminMarker;
578 use fidl_fuchsia_fs_startup::{
579 CreateOptions, MountOptions, StartOptions, StartupMarker, VolumesMarker,
580 };
581 use fidl_fuchsia_fxfs::DebugMarker;
582 use fidl_fuchsia_io as fio;
583 use fidl_fuchsia_process_lifecycle::{LifecycleMarker, LifecycleProxy};
584 use fuchsia_async as fasync;
585 use fuchsia_component_client::connect_to_protocol_at_dir_svc;
586 use fuchsia_fs::directory::readdir;
587 use futures::future::{BoxFuture, FusedFuture, FutureExt};
588 use futures::{pin_mut, select};
589 use fxfs::filesystem::FxFilesystem;
590 use fxfs::object_store::volume::root_volume;
591 use fxfs::object_store::{NewChildStoreOptions, StoreOptions};
592 use std::pin::Pin;
593 use std::sync::Arc;
594 use storage_device::DeviceHolder;
595 use storage_device::block_device::BlockDevice;
596 use test_vmo_backed_block_server::VmoBackedServer;
597
598 async fn run_test(
599 callback: impl Fn(&fio::DirectoryProxy, LifecycleProxy) -> BoxFuture<'static, ()>,
600 ) -> (Pin<Box<impl FusedFuture>>, Arc<VmoBackedServer>) {
601 const BLOCK_SIZE: u32 = 512;
602 let block_server = Arc::new(
603 VmoBackedServer::new(16384, BLOCK_SIZE, &[]).expect("Failed to create VmoBackedServer"),
604 );
605
606 {
607 let fs = FxFilesystem::new_empty(DeviceHolder::new(
608 BlockDevice::new(
609 new_block_client(block_server.connect())
610 .await
611 .expect("Unable to create block client"),
612 false,
613 )
614 .await
615 .unwrap(),
616 ))
617 .await
618 .expect("FxFilesystem::new_empty failed");
619 {
620 let root_volume = root_volume(fs.clone()).await.expect("Open root_volume failed");
621 root_volume
622 .new_volume("default", NewChildStoreOptions::default())
623 .await
624 .expect("Create volume failed");
625 }
626 fs.close().await.expect("close failed");
627 }
628
629 let (client_end, server_end) = fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
630
631 let (lifecycle_client, lifecycle_server) =
632 fidl::endpoints::create_proxy::<LifecycleMarker>();
633
634 let mut component_task = Box::pin(
635 async {
636 Component::new()
637 .run(server_end.into_channel(), Some(lifecycle_server.into_channel()))
638 .await
639 .expect("Failed to run component");
640 }
641 .fuse(),
642 );
643
644 let startup_proxy = connect_to_protocol_at_dir_svc::<StartupMarker>(&client_end)
645 .expect("Unable to connect to Startup protocol");
646 let block_server_connection = block_server.connect();
647 let task = async {
648 startup_proxy
649 .start(block_server_connection, &StartOptions::default())
650 .await
651 .expect("Start failed (FIDL)")
652 .expect("Start failed");
653 callback(&client_end, lifecycle_client).await;
654 }
655 .fuse();
656
657 pin_mut!(task);
658
659 loop {
660 select! {
661 () = component_task => {},
662 () = task => break,
663 }
664 }
665
666 (component_task, block_server)
667 }
668
669 #[fuchsia::test(threads = 2)]
670 async fn test_clear_caches() {
671 let (component_task, _block_server) = run_test(|client, _| {
672 let debug_proxy = connect_to_protocol_at_dir_svc::<DebugMarker>(client)
673 .expect("Unable to connect to Debug protocol");
674 let admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(client)
675 .expect("Unable to connect to Admin protocol");
676 async move {
677 debug_proxy
678 .clear_caches()
679 .await
680 .expect("clear_caches failed (FIDL)")
681 .expect("clear_caches failed");
682 admin_proxy.shutdown().await.expect("shutdown failed");
683 }
684 .boxed()
685 })
686 .await;
687 component_task.await;
688 }
689
690 #[fuchsia::test(threads = 2)]
691 async fn test_shutdown() {
692 let (component_task, _block_server) = run_test(|client, _| {
693 let admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(client)
694 .expect("Unable to connect to Admin protocol");
695 async move {
696 admin_proxy.shutdown().await.expect("shutdown failed");
697 }
698 .boxed()
699 })
700 .await;
701 assert!(!component_task.is_terminated());
702 }
703
704 #[fuchsia::test(threads = 2)]
705 async fn test_lifecycle_stop() {
706 let (component_task, _block_server) = run_test(|_, lifecycle_client| {
707 lifecycle_client.stop().expect("Stop failed");
708 async move {
709 fasync::OnSignals::new(
710 &lifecycle_client.into_channel().expect("into_channel failed"),
711 zx::Signals::CHANNEL_PEER_CLOSED,
712 )
713 .await
714 .expect("OnSignals failed");
715 }
716 .boxed()
717 })
718 .await;
719 component_task.await;
720 }
721
722 #[fuchsia::test(threads = 2)]
723 async fn test_create_and_remove() {
724 let (component_task, _block_server) = run_test(|client, _| {
725 let volumes_proxy = connect_to_protocol_at_dir_svc::<VolumesMarker>(client)
726 .expect("Unable to connect to Volumes protocol");
727
728 let fs_admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(client)
729 .expect("Unable to connect to Admin protocol");
730
731 async move {
732 let (dir_proxy, server_end) =
733 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
734 volumes_proxy
735 .create("test", server_end, CreateOptions::default(), MountOptions::default())
736 .await
737 .expect("fidl failed")
738 .expect("create failed");
739
740 volumes_proxy
742 .remove("test")
743 .await
744 .expect("fidl failed")
745 .expect_err("remove succeeded");
746
747 let volume_admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(&dir_proxy)
748 .expect("Unable to connect to Admin protocol");
749 volume_admin_proxy.shutdown().await.expect("shutdown failed");
750
751 let (_dir_proxy, server_end) =
753 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
754 volumes_proxy
755 .create("test", server_end, CreateOptions::default(), MountOptions::default())
756 .await
757 .expect("fidl failed")
758 .expect_err("create succeeded");
759
760 volumes_proxy.remove("test").await.expect("fidl failed").expect("remove failed");
761
762 volumes_proxy
764 .remove("test")
765 .await
766 .expect("fidl failed")
767 .expect_err("remove failed");
768
769 let (_dir_proxy, server_end) =
771 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
772 volumes_proxy
773 .create("test", server_end, CreateOptions::default(), MountOptions::default())
774 .await
775 .expect("fidl failed")
776 .expect("create failed");
777
778 fs_admin_proxy.shutdown().await.expect("shutdown failed");
779 }
780 .boxed()
781 })
782 .await;
783 component_task.await;
784 }
785
786 #[fuchsia::test(threads = 2)]
787 async fn test_volumes_enumeration() {
788 let (component_task, _block_server) = run_test(|client, _| {
789 let volumes_proxy = connect_to_protocol_at_dir_svc::<VolumesMarker>(client)
790 .expect("Unable to connect to Volumes protocol");
791
792 let (volumes_dir_proxy, server_end) =
793 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
794 client
795 .open("volumes", fio::PERM_READABLE, &Default::default(), server_end.into_channel())
796 .expect("open failed");
797
798 let fs_admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(client)
799 .expect("Unable to connect to Admin protocol");
800
801 async move {
802 let (_dir_proxy, server_end) =
803 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
804 volumes_proxy
805 .create("test", server_end, CreateOptions::default(), MountOptions::default())
806 .await
807 .expect("fidl failed")
808 .expect("create failed");
809
810 let entries = readdir(&volumes_dir_proxy).await.expect("readdir failed");
811 let mut entry_names = entries.iter().map(|d| d.name.as_str()).collect::<Vec<_>>();
812 entry_names.sort();
813 assert_eq!(entry_names, ["default", "test"]);
814
815 fs_admin_proxy.shutdown().await.expect("shutdown failed");
816 }
817 .boxed()
818 })
819 .await;
820 component_task.await;
821 }
822
823 #[fuchsia::test(threads = 2)]
824 async fn test_create_with_guid() {
825 let (component_task, block_server) = run_test(|client, _| {
826 let volumes_proxy = connect_to_protocol_at_dir_svc::<VolumesMarker>(client)
827 .expect("Unable to connect to Volumes protocol");
828 let admin_proxy = connect_to_protocol_at_dir_svc::<AdminMarker>(client)
829 .expect("Unable to connect to Admin protocol");
830 async move {
831 let (_dir_proxy, server_end) =
832 fidl::endpoints::create_proxy::<fio::DirectoryMarker>();
833 let guid = [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16];
834 volumes_proxy
835 .create(
836 "test_guid",
837 server_end,
838 CreateOptions { guid: Some(guid), ..Default::default() },
839 MountOptions::default(),
840 )
841 .await
842 .expect("fidl failed")
843 .expect("create failed");
844 admin_proxy.shutdown().await.expect("shutdown failed");
845 }
846 .boxed()
847 })
848 .await;
849 component_task.await;
850
851 let fs = FxFilesystem::open(DeviceHolder::new(
853 BlockDevice::new(
854 new_block_client(block_server.connect())
855 .await
856 .expect("Unable to create block client"),
857 false,
858 )
859 .await
860 .unwrap(),
861 ))
862 .await
863 .expect("FxFilesystem::open failed");
864 {
865 let root_volume = root_volume(fs.clone()).await.expect("Open root_volume failed");
866 let vol = root_volume
867 .volume("test_guid", StoreOptions::default())
868 .await
869 .expect("Open volume failed");
870 assert_eq!(vol.guid(), [1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16]);
871 }
872 fs.close().await.expect("close failed");
873 }
874}