gpt_component/
partitions_directory.rs1use 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
17pub 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 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 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
67pub 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}