1use super::{Operation, RequestId, TraceFlowId};
6use crate::{IntoOrchestrator, callback_interface};
7use fidl::endpoints::RequestStream;
8use fidl_fuchsia_storage_block as fblock;
9use fidl_fuchsia_storage_block::MAX_TRANSFER_UNBOUNDED;
10use fuchsia_async as fasync;
11use fuchsia_sync::{Condvar, Mutex};
12use futures::stream::{AbortHandle, Abortable};
13use std::borrow::{Borrow, Cow};
14use std::ffi::{CStr, c_char, c_void};
15use std::num::NonZero;
16use std::sync::Arc;
17
18pub type Session = callback_interface::Session<InterfaceAdapter>;
20
21#[repr(C)]
22pub struct Callbacks {
23 pub context: *mut c_void,
27 pub start_thread: unsafe extern "C" fn(context: *mut c_void, arg: *const c_void),
33 pub on_new_session: unsafe extern "C" fn(context: *mut c_void, session: *const Session),
38 pub on_requests:
45 unsafe extern "C" fn(context: *mut c_void, requests: *mut Request, request_count: usize),
46 pub log: unsafe extern "C" fn(context: *mut c_void, message: *const c_char, message_len: usize),
49}
50
51impl Callbacks {
52 #[allow(dead_code)]
53 fn log(&self, msg: &str) {
54 let msg = msg.as_bytes();
55 unsafe {
57 (self.log)(self.context, msg.as_ptr() as *const c_char, msg.len());
58 }
59 }
60}
61
62#[allow(dead_code)]
64pub struct UnownedVmo(zx::sys::zx_handle_t);
65
66#[repr(C)]
67pub struct Request {
68 pub request_id: RequestId,
69 pub operation: Operation,
70 pub trace_flow_id: TraceFlowId,
71 pub vmo: UnownedVmo,
72}
73
74unsafe impl Send for Callbacks {}
75unsafe impl Sync for Callbacks {}
76
77pub struct InterfaceAdapter {
79 callbacks: Callbacks,
80 info: super::DeviceInfo,
81}
82
83impl callback_interface::Interface for InterfaceAdapter {
84 type Orchestrator = Orchestrator;
85
86 fn get_info(&self) -> Cow<'_, super::DeviceInfo> {
87 Cow::Borrowed(&self.info)
88 }
89
90 fn spawn_session(&self, session: Arc<Session>) {
91 unsafe {
92 (self.callbacks.on_new_session)(self.callbacks.context, Arc::into_raw(session));
93 }
94 }
95
96 fn on_requests(&self, requests: &[callback_interface::Request]) {
97 let mut c_requests = Vec::with_capacity(requests.len());
98 for req in requests {
99 c_requests.push(Request {
100 request_id: req.request_id,
101 operation: req.operation.clone(),
102 trace_flow_id: req.trace_flow_id,
103 vmo: UnownedVmo(
107 req.vmo.as_ref().map(|v| v.raw_handle()).unwrap_or(zx::sys::ZX_HANDLE_INVALID),
108 ),
109 });
110 }
111 unsafe {
112 (self.callbacks.on_requests)(
113 self.callbacks.context,
114 c_requests.as_mut_ptr(),
115 c_requests.len(),
116 )
117 }
118 }
119}
120
121#[repr(C)]
122pub struct PartitionInfo {
123 pub device_flags: u32,
124 pub start_block: u64,
125 pub block_count: u64,
126 pub block_size: u32,
127 pub type_guid: [u8; 16],
128 pub instance_guid: [u8; 16],
129 pub name: *const c_char,
130 pub flags: u64,
131 pub max_transfer_size: u32,
132}
133
134#[allow(non_camel_case_types)]
136type zx_handle_t = zx::sys::zx_handle_t;
137
138#[allow(non_camel_case_types)]
140type zx_status_t = zx::sys::zx_status_t;
141
142impl PartitionInfo {
143 unsafe fn to_rust(&self) -> super::DeviceInfo {
147 super::DeviceInfo::Partition(super::PartitionInfo {
148 device_flags: fblock::DeviceFlag::from_bits_truncate(self.device_flags),
149 start_block_offset: Some(self.start_block),
150 block_count: self.block_count,
151 type_guid: self.type_guid,
152 instance_guid: self.instance_guid,
153 name: if self.name.is_null() {
154 "".to_string()
155 } else {
156 String::from_utf8_lossy(unsafe { CStr::from_ptr(self.name).to_bytes() }).to_string()
157 },
158 flags: Some(self.flags),
159 max_transfer_blocks: if self.max_transfer_size != MAX_TRANSFER_UNBOUNDED {
160 NonZero::new(self.max_transfer_size / self.block_size)
161 } else {
162 None
163 },
164 })
165 }
166}
167
168struct ExecutorMailbox(Mutex<Mail>, Condvar);
169
170impl ExecutorMailbox {
171 fn post(&self, mail: Mail) -> Mail {
172 let old = std::mem::replace(&mut *self.0.lock(), mail);
173 self.1.notify_all();
174 old
175 }
176
177 fn new() -> Self {
178 Self(Mutex::default(), Condvar::new())
179 }
180}
181
182type ShutdownCallback = unsafe extern "C" fn(*mut c_void);
183
184#[derive(Default)]
185enum Mail {
186 #[default]
187 None,
188 Initialized(fasync::ScopeHandle, AbortHandle),
189 AsyncShutdown(*const BlockServer, ShutdownCallback, *mut c_void),
190 ThreadFinished(*const BlockServer, ShutdownCallback, *mut c_void),
191 Finished,
192}
193
194unsafe impl Send for Mail {}
196
197pub struct Orchestrator {
198 session_manager: callback_interface::SessionManager<InterfaceAdapter>,
199 mbox: ExecutorMailbox,
200}
201
202impl IntoOrchestrator for Arc<Orchestrator> {
203 type SM = callback_interface::SessionManager<InterfaceAdapter>;
204
205 fn into_orchestrator(self) -> Arc<Orchestrator> {
206 self
207 }
208}
209
210impl Borrow<callback_interface::SessionManager<InterfaceAdapter>> for Orchestrator {
211 fn borrow(&self) -> &callback_interface::SessionManager<InterfaceAdapter> {
212 &self.session_manager
213 }
214}
215
216pub struct BlockServer {
217 server: super::BlockServer<callback_interface::SessionManager<InterfaceAdapter>>,
218 scope: fasync::ScopeHandle,
219 abort_handle: AbortHandle,
220 orchestrator: Arc<Orchestrator>,
221}
222
223#[unsafe(no_mangle)]
230pub unsafe extern "C" fn block_server_new(
231 partition_info: &PartitionInfo,
232 callbacks: Callbacks,
233) -> *mut BlockServer {
234 let start_thread = callbacks.start_thread;
235 let context = callbacks.context;
236
237 let session_manager = callback_interface::SessionManager::new(Arc::new(InterfaceAdapter {
238 callbacks,
239 info: unsafe { partition_info.to_rust() },
240 }));
241
242 let orchestrator = Arc::new(Orchestrator { session_manager, mbox: ExecutorMailbox::new() });
243
244 unsafe {
245 (start_thread)(context, Arc::into_raw(orchestrator.clone()) as *const c_void);
246 }
247
248 let mbox = &orchestrator.mbox;
249 let mail = {
250 let mut mail = mbox.0.lock();
251 mbox.1.wait_while(&mut mail, |mail| matches!(mail, Mail::None));
252 std::mem::replace(&mut *mail, Mail::None)
253 };
254
255 let block_size = partition_info.block_size;
256 match mail {
257 Mail::Initialized(scope, abort_handle) => Box::into_raw(Box::new(BlockServer {
258 server: super::BlockServer::new(block_size, orchestrator.clone()),
259 scope,
260 abort_handle,
261 orchestrator: orchestrator.clone(),
262 })),
263 Mail::Finished => std::ptr::null_mut(),
264 _ => unreachable!(),
265 }
266}
267
268#[unsafe(no_mangle)]
277pub unsafe extern "C" fn block_server_thread(arg: *const c_void) {
278 let orchestrator = unsafe { &*(arg as *const Orchestrator) };
279
280 let mut executor = fasync::LocalExecutor::default();
281 let scope = fasync::Scope::new();
282
283 let (abort_handle, registration) = AbortHandle::new_pair();
286 let root_task = scope.spawn(async move {
287 let _ = Abortable::new(std::future::pending::<()>(), registration).await;
288 });
289 orchestrator.mbox.post(Mail::Initialized(scope.clone(), abort_handle));
290
291 let _ = executor.run_singlethreaded(root_task);
294
295 {
297 let mut mbox = orchestrator.mbox.0.lock();
298 let mail = std::mem::take(&mut *mbox);
299 if let Mail::AsyncShutdown(block_server, callback, arg) = mail {
300 *mbox = Mail::ThreadFinished(block_server, callback, arg);
301 orchestrator.mbox.1.notify_all();
302 } else {
303 *mbox = mail;
304 }
305 }
306
307 let _ = executor.run_singlethreaded(scope.cancel());
309
310 orchestrator.session_manager.terminate();
314}
315
316#[unsafe(no_mangle)]
323pub unsafe extern "C" fn block_server_thread_release(arg: *const c_void) {
324 let orchestrator = unsafe { Arc::from_raw(arg as *const Orchestrator) };
326
327 let mail = orchestrator.mbox.post(Mail::Finished);
328 match mail {
329 Mail::None | Mail::Finished => {}
330 Mail::ThreadFinished(block_server, callback, arg) => {
331 let _ = unsafe { Box::from_raw(block_server as *mut BlockServer) };
334
335 unsafe {
337 callback(arg);
338 }
339 }
340 _ => panic!("block_server_thread_release called while thread is still running"),
341 }
342}
343
344#[unsafe(no_mangle)]
349pub unsafe extern "C" fn block_server_delete(block_server: *const BlockServer) {
350 {
351 let server = unsafe { &*block_server };
353
354 server.abort_handle.abort();
360 {
361 let mbox = &server.orchestrator.mbox;
362 let mut mail = mbox.0.lock();
363 mbox.1.wait_while(&mut mail, |mbox| !matches!(mbox, Mail::Finished));
364 }
365
366 Borrow::<callback_interface::SessionManager<InterfaceAdapter>>::borrow(
368 server.orchestrator.as_ref(),
369 )
370 .terminate();
371 }
372
373 let _ = unsafe { Box::from_raw(block_server as *mut BlockServer) };
375}
376
377#[unsafe(no_mangle)]
382pub unsafe extern "C" fn block_server_delete_async(
383 block_server: *const BlockServer,
384 callback: ShutdownCallback,
385 arg: *mut c_void,
386) {
387 let abort_handle = {
388 let server = unsafe { &*block_server };
390
391 assert!(!matches!(
394 server.orchestrator.mbox.post(Mail::AsyncShutdown(block_server, callback, arg)),
395 Mail::Finished
396 ));
397
398 server.abort_handle.clone()
399 };
400
401 abort_handle.abort();
403}
404
405#[unsafe(no_mangle)]
411pub unsafe extern "C" fn block_server_serve(block_server: *const BlockServer, handle: zx_handle_t) {
412 let block_server = unsafe { &*block_server };
413 let handle = unsafe { zx::NullableHandle::from_raw(handle) };
414 block_server.scope.spawn(async move {
415 let _ = block_server
416 .server
417 .handle_requests(fblock::BlockRequestStream::from_channel(
418 fasync::Channel::from_channel(handle.into()),
419 ))
420 .await;
421 });
422}
423
424#[unsafe(no_mangle)]
428pub unsafe extern "C" fn block_server_session_run(session: &Session) {
429 let session = unsafe { Arc::from_raw(session) };
430 session.run();
431 let _ = Arc::into_raw(session);
432}
433
434#[unsafe(no_mangle)]
438pub unsafe extern "C" fn block_server_session_release(session: &Session) {
439 session.terminate_async();
440 unsafe { Arc::from_raw(session) };
441}
442
443#[unsafe(no_mangle)]
447pub unsafe extern "C" fn block_server_send_reply(
448 block_server: &BlockServer,
449 request_id: RequestId,
450 status: zx_status_t,
451) {
452 block_server
453 .orchestrator
454 .session_manager
455 .complete_request(request_id, zx::Status::from_raw(status));
456}