Skip to main content

fxfs_platform_testing/fuchsia/
component.rs

1// Copyright 2022 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use crate::fuchsia::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
57/// Runs Fxfs as a component.
58pub struct Component {
59    state: futures::lock::Mutex<State>,
60
61    // The execution scope of the pseudo filesystem.
62    scope: ExecutionScope,
63
64    // The root of the pseudo filesystem for the component.
65    export_dir: Arc<vfs::directory::immutable::Simple>,
66}
67
68// Wrapper type to add `FsInspect` support to an `OpenFxFilesystem`.
69struct InspectedFxFilesystem(OpenFxFilesystem, /*fs_id=*/ 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            // TODO(https://fxbug.dev/42175930): Should these be moved to per-volume nodes?
106            total_nodes: 0,
107            used_nodes: 0,
108        }
109    }
110}
111
112struct RunningState {
113    // We have to wrap this in an Arc, even though it itself basically just wraps an Arc, so that
114    // FsInspectTree can reference `fs` as a Weak<dyn FsInspect>`.
115    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", /* must_be_directory: */ false);
132            let _ = outgoing_dir.remove_entry("debug", /* must_be_directory: */ 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    /// Runs Fxfs as a component.
149    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        // TODO(b/315704445): Only enable in debug builds?
211        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        // TODO(https://fxbug.dev/42174810): This is not very graceful.  It would be better for the client to
291        // explicitly shut down all volumes first, and make this fail if there are remaining active
292        // connections.  Fix the bug in fs_test which requires this.
293        state.stop(&self.export_dir).await;
294        let client = new_block_client(device).await?;
295
296        // TODO(https://fxbug.dev/42063349) Add support for block sizes greater than the page size.
297        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            /* overwrite: */ 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            // Unwrap ok, shouldn't have anything else recording or replaying this early in startup.
365            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?, /* read_only: */ 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, /* read_only: */ 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        // TODO(b/311550633): Stash ok res in inspect.
408        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    // Returns true if we should close the connection.
422    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    /// Handles fuchsia.fxfs.Debug requests, providing live debugging internals of the running
440    /// filesystem.
441    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                // This should fail whilst the volume is mounted.
741                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                // Creating another volume with the same name should fail.
752                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                // Removing a non-existent volume should fail.
763                volumes_proxy
764                    .remove("test")
765                    .await
766                    .expect("fidl failed")
767                    .expect_err("remove failed");
768
769                // Create the same volume again and it should now succeed.
770                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        // Verify the GUID
852        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}