Skip to main content

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