Skip to main content

gpt_component/
partitions_directory.rs

1// Copyright 2024 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::gpt::GptManager;
6use block_server::{BlockServer, SessionManager};
7use fidl::endpoints::RequestStream as _;
8use fidl_fuchsia_storage_block as fblock;
9use fidl_fuchsia_storage_partitions as fpartitions;
10use fuchsia_async as fasync;
11use fuchsia_sync::Mutex;
12use futures::future::join_all;
13use std::collections::BTreeMap;
14use std::sync::{Arc, Weak};
15use vfs::directory::helper::DirectlyMutable as _;
16
17/// A directory of instances of the fuchsia.storage.partitions.PartitionService service.
18pub struct PartitionsDirectory {
19    node: Arc<vfs::directory::immutable::Simple>,
20    entries: Mutex<BTreeMap<String, PartitionsDirectoryEntry>>,
21}
22
23impl std::fmt::Debug for PartitionsDirectory {
24    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> {
25        f.debug_struct("PartitionDirectory").field("entries", &self.entries).finish()
26    }
27}
28
29impl PartitionsDirectory {
30    pub fn new(node: Arc<vfs::directory::immutable::Simple>) -> Self {
31        Self { node, entries: Default::default() }
32    }
33
34    pub async fn clear(&self) {
35        self.node.remove_all_entries();
36        let entries = std::mem::take(&mut *self.entries.lock());
37        join_all(entries.into_values().map(|entry| entry.scope.cancel())).await;
38    }
39
40    /// Adds an entry for a GPT partition.  Serves the "volume" and "partition" protocols.
41    pub fn add_partition<SM: SessionManager + Send + Sync + 'static>(
42        &self,
43        name: &str,
44        block_server: Weak<BlockServer<SM>>,
45        gpt_manager: Weak<GptManager>,
46        gpt_index: usize,
47    ) {
48        let entry = PartitionsDirectoryEntry::new_partition(block_server, gpt_manager, gpt_index);
49        self.node.add_entry(name, entry.node.clone()).expect("Added an entry twice");
50        self.entries.lock().insert(name.to_string(), entry);
51    }
52
53    /// Adds an entry for a composite partition.  Serves the "volume" and "overlay" protocols.
54    pub fn add_composite<SM: SessionManager + Send + Sync + 'static>(
55        &self,
56        name: &str,
57        block_server: Weak<BlockServer<SM>>,
58        gpt_manager: Weak<GptManager>,
59        gpt_indexes: Vec<usize>,
60    ) {
61        let entry = PartitionsDirectoryEntry::new_composite(block_server, gpt_manager, gpt_indexes);
62        self.node.add_entry(name, entry.node.clone()).expect("Added an entry twice");
63        self.entries.lock().insert(name.to_string(), entry);
64    }
65}
66
67/// A node which hosts an instance of fuchsia.storage.partitions.PartitionService.
68pub struct PartitionsDirectoryEntry {
69    scope: fasync::Scope,
70    node: Arc<vfs::directory::immutable::Simple>,
71}
72
73impl std::fmt::Debug for PartitionsDirectoryEntry {
74    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> Result<(), std::fmt::Error> {
75        f.debug_struct("PartitionDirectoryEntry").finish()
76    }
77}
78
79impl PartitionsDirectoryEntry {
80    fn new_partition<SM: SessionManager + Send + Sync + 'static>(
81        block_server: Weak<BlockServer<SM>>,
82        gpt_manager: Weak<GptManager>,
83        gpt_index: usize,
84    ) -> Self {
85        let scope = fasync::Scope::new();
86        let has_mapper = gpt_manager.upgrade().is_some_and(|m| m.has_mapper());
87        let node = vfs::directory::immutable::simple();
88        let volume_server = block_server.clone();
89        node.add_entry(
90            "volume",
91            vfs::service::endpoint({
92                let scope = scope.to_handle();
93                move |_scope, channel| {
94                    let server = volume_server.clone();
95                    let requests = fblock::BlockRequestStream::from_channel(channel);
96                    scope.spawn(async move {
97                        if let Some(server) = server.upgrade() {
98                            if let Err(err) = server.handle_requests(requests).await {
99                                log::error!(err:?; "Error handling requests");
100                            }
101                        }
102                    });
103                }
104            }),
105        )
106        .unwrap();
107        node.add_entry(
108            "partition",
109            vfs::service::endpoint({
110                let scope = scope.to_handle();
111                move |_scope, channel| {
112                    let manager = gpt_manager.clone();
113                    let requests = fpartitions::PartitionRequestStream::from_channel(channel);
114                    scope.spawn(async move {
115                        if let Some(manager) = manager.upgrade() {
116                            if let Err(err) =
117                                manager.handle_partitions_requests(gpt_index, requests).await
118                            {
119                                log::error!(err:?; "Error handling requests");
120                            }
121                        }
122                    });
123                }
124            }),
125        )
126        .unwrap();
127        if has_mapper {
128            let mapper_server = block_server;
129            node.add_entry(
130                "mapper",
131                vfs::service::endpoint({
132                    let scope = scope.to_handle();
133                    move |_scope, channel| {
134                        let server = mapper_server.clone();
135                        let requests = fblock::MapperRequestStream::from_channel(channel);
136                        scope.spawn(async move {
137                            if let Some(server) = server.upgrade() {
138                                if let Err(err) = server.handle_mapper_requests(requests).await {
139                                    log::error!(err:?; "Error handling mapper requests");
140                                }
141                            }
142                        });
143                    }
144                }),
145            )
146            .unwrap();
147        }
148
149        Self { scope, node }
150    }
151
152    fn new_composite<SM: SessionManager + Send + Sync + 'static>(
153        block_server: Weak<BlockServer<SM>>,
154        gpt_manager: Weak<GptManager>,
155        gpt_indexes: Vec<usize>,
156    ) -> Self {
157        let scope = fasync::Scope::new();
158        let has_mapper = gpt_manager.upgrade().is_some_and(|m| m.has_mapper());
159        let node = vfs::directory::immutable::simple();
160        let volume_server = block_server.clone();
161        node.add_entry(
162            "volume",
163            vfs::service::endpoint({
164                let scope = scope.to_handle();
165                move |_scope, channel| {
166                    let server = volume_server.clone();
167                    let requests = fblock::BlockRequestStream::from_channel(channel);
168                    scope.spawn(async move {
169                        if let Some(server) = server.upgrade() {
170                            if let Err(err) = server.handle_requests(requests).await {
171                                log::error!(err:?; "Error handling requests");
172                            }
173                        }
174                    });
175                }
176            }),
177        )
178        .unwrap();
179        node.add_entry(
180            "overlay",
181            vfs::service::endpoint({
182                let scope = scope.to_handle();
183                move |_scope, channel| {
184                    let manager = gpt_manager.clone();
185                    let gpt_indexes = gpt_indexes.clone();
186                    let requests =
187                        fpartitions::OverlayPartitionRequestStream::from_channel(channel);
188                    scope.spawn(async move {
189                        if let Some(manager) = manager.upgrade() {
190                            if let Err(err) = manager
191                                .handle_composite_partitions_requests(gpt_indexes, requests)
192                                .await
193                            {
194                                log::error!(err:?; "Error handling requests");
195                            }
196                        }
197                    });
198                }
199            }),
200        )
201        .unwrap();
202        if has_mapper {
203            let mapper_server = block_server;
204            node.add_entry(
205                "mapper",
206                vfs::service::endpoint({
207                    let scope = scope.to_handle();
208                    move |_scope, channel| {
209                        let server = mapper_server.clone();
210                        let requests = fblock::MapperRequestStream::from_channel(channel);
211                        scope.spawn(async move {
212                            if let Some(server) = server.upgrade() {
213                                if let Err(err) = server.handle_mapper_requests(requests).await {
214                                    log::error!(err:?; "Error handling mapper requests");
215                                }
216                            }
217                        });
218                    }
219                }),
220            )
221            .unwrap();
222        }
223
224        Self { scope, node }
225    }
226}