1use 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
61pub struct Component {
63 state: futures::lock::Mutex<State>,
64
65 scope: ExecutionScope,
67
68 export_dir: Arc<vfs::directory::immutable::Simple>,
70}
71
72struct InspectedFxFilesystem(OpenFxFilesystem, 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 total_nodes: 0,
111 used_nodes: 0,
112 }
113 }
114}
115
116struct RunningState {
117 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", false);
136 let _ = outgoing_dir.remove_entry("debug", 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 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 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 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 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 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 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?, 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, 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 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 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 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 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 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 volumes_proxy
795 .remove("test")
796 .await
797 .expect("fidl failed")
798 .expect_err("remove failed");
799
800 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 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}