1use crate::verifier::Verifier;
6use crate::{
7 ActiveRequests, DecodedRequest, DeviceInfo, HandleRequestResult, IntoOrchestrator, OffsetMap,
8 Operation, RequestId, SessionHelper, TraceFlowId, WriteFlags,
9};
10use anyhow::Error;
11use block_protocol::{BlockFifoRequest, BlockFifoResponse};
12use fidl_fuchsia_storage_block as fblock;
13use fuchsia_sync::{Condvar, Mutex};
14use futures::TryStreamExt as _;
15use futures::stream::{AbortHandle, Abortable};
16use mapping::reader::BlockService;
17use std::borrow::{Borrow, Cow};
18use std::collections::{HashMap, VecDeque};
19use std::mem::MaybeUninit;
20use std::sync::{Arc, OnceLock, Weak};
21
22pub mod block_service;
23pub use block_service::DefaultCallbackBlockService;
24
25#[derive(Clone, Debug)]
27pub struct Request {
28 pub request_id: RequestId,
33 pub operation: Operation,
34 pub trace_flow_id: TraceFlowId,
35 pub vmo: Option<Arc<zx::Vmo>>,
37}
38
39pub trait Interface: Send + Sync + Unpin + 'static {
40 type Orchestrator: Borrow<SessionManager<Self>> + Send + Sync;
41
42 fn get_info(&self) -> Cow<'_, DeviceInfo>;
44
45 fn spawn_session(&self, session: Arc<Session<Self>>);
48
49 fn on_requests(&self, requests: &[Request]);
56
57 fn into_block_service(
61 self: Arc<Self>,
62 orchestrator: &Arc<Self::Orchestrator>,
63 ) -> Arc<dyn BlockService> {
64 Arc::new(DefaultCallbackBlockService::<Self>::new(orchestrator))
65 }
66
67 fn on_open_mapper_session(
70 &self,
71 _mapping_vmo: &zx::Vmo,
72 delivery_queue: zx::Vmo,
73 ) -> Result<Verifier, zx::Status> {
74 Ok(Verifier::new(delivery_queue))
75 }
76}
77
78const INTERNAL_REQUEST_FLAG: usize = 1 << 63;
83
84#[derive(Default)]
85struct InflightRequests {
86 count: usize,
88 callbacks: HashMap<RequestId, Box<dyn FnOnce(Result<(), zx::Status>) + Send>>,
90 next_internal_id: usize,
92}
93
94struct SessionManagerInner<I: Interface + ?Sized> {
95 open_sessions: HashMap<usize, Weak<Session<I>>>,
96}
97
98const FIFO_WAKE_SIGNAL: zx::Signals = zx::Signals::USER_0;
102
103const SHUTDOWN_SIGNAL: zx::Signals = zx::Signals::USER_1;
105
106pub struct SessionManager<I: Interface + ?Sized> {
107 interface: Arc<I>,
108 block_size: u32,
109 active_requests: ActiveRequests<Arc<Session<I>>>,
111 inflight_requests: Mutex<InflightRequests>,
112 no_inflight_requests_condvar: Condvar,
113 inner: Mutex<SessionManagerInner<I>>,
114 no_open_sessions_condvar: Condvar,
115 block_service: OnceLock<Arc<dyn BlockService>>,
116}
117
118impl<I: Interface + ?Sized> super::SessionManager for SessionManager<I> {
119 const SUPPORTS_DECOMPRESSION: bool = false;
120
121 type Orchestrator = I::Orchestrator;
122 type Session = Arc<Session<I>>;
123
124 fn session_eq(a: &Arc<Session<I>>, b: &Arc<Session<I>>) -> bool {
125 Arc::ptr_eq(a, b)
126 }
127
128 async fn on_attach_vmo(
129 _orchestrator: Arc<Self::Orchestrator>,
130 _vmo: &Arc<zx::Vmo>,
131 ) -> Result<(), zx::Status> {
132 Ok(())
133 }
134
135 async fn open_session(
140 orchestrator: Arc<Self::Orchestrator>,
141 mut stream: fblock::SessionRequestStream,
142 offset_map: OffsetMap,
143 block_size: u32,
144 ) -> Result<(), Error> {
145 let sm: &SessionManager<I> = orchestrator.as_ref().borrow();
146 let max_blocks = sm.get_info().max_transfer_blocks();
147 let (helper, fifo) =
148 SessionHelper::new(orchestrator.clone(), offset_map, max_blocks, block_size)?;
149 let (abort_handle, registration) = AbortHandle::new_pair();
150 let session = Arc::new(Session {
151 helper,
152 fifo,
153 queue: Mutex::default(),
154 abort_handle,
155 close_callback: Mutex::new(None),
156 });
157 let sm = orchestrator.as_ref().borrow();
158 sm.inner
159 .lock()
160 .open_sessions
161 .insert(Arc::as_ptr(&session) as usize, Arc::downgrade(&session));
162
163 sm.interface.spawn_session(session.clone());
164
165 let result = Abortable::new(
166 async {
167 while let Some(request) = stream.try_next().await? {
168 match session.helper.handle_request(request).await? {
169 HandleRequestResult::Ok => {}
170 HandleRequestResult::Closed(callback) => {
171 *session.close_callback.lock() = Some(callback);
172 break;
173 }
174 }
175 }
176 Ok(())
177 },
178 registration,
179 )
180 .await
181 .unwrap_or_else(|e| Err(e.into()));
182
183 let _ = session.fifo.signal(zx::Signals::empty(), SHUTDOWN_SIGNAL);
184
185 result
186 }
187
188 fn open_mapper_session(
189 orchestrator: Arc<Self::Orchestrator>,
190 session: fidl::endpoints::ServerEnd<fblock::MapperSessionMarker>,
191 mapping_vmo: zx::Vmo,
192 _block_size: u32,
193 port: Option<zx::Port>,
194 delivery_queue: Option<zx::Vmo>,
195 ) -> Result<impl Future<Output = Result<(), Error>> + Send, zx::Status> {
196 let sm: &SessionManager<I> = orchestrator.as_ref().borrow();
197 let service = sm.into_block_service(&orchestrator);
198 let interface = sm.interface.clone();
199 crate::mapper::serve_mapper_session(
200 Arc::new(move |mapping_vmo: &zx::Vmo, dq: zx::Vmo| {
201 interface.on_open_mapper_session(mapping_vmo, dq)
202 }),
203 service,
204 session,
205 mapping_vmo,
206 port,
207 delivery_queue,
208 )
209 }
210
211 fn get_info(&self) -> Cow<'_, super::DeviceInfo> {
212 self.interface.get_info()
213 }
214
215 fn active_requests(&self) -> &ActiveRequests<Arc<Session<I>>> {
216 &self.active_requests
217 }
218}
219
220impl<I: Interface + ?Sized> SessionManager<I> {
221 pub fn into_block_service(&self, orchestrator: &Arc<I::Orchestrator>) -> Arc<dyn BlockService> {
223 self.block_service
224 .get_or_init(|| self.interface.clone().into_block_service(orchestrator))
225 .clone()
226 }
227}
228
229impl<I: Interface + ?Sized> SessionManager<I> {
230 pub fn new(interface: Arc<I>, block_size: u32) -> Self
231 where
232 I: Sized,
233 {
234 Self {
235 interface,
236 block_size,
237 active_requests: ActiveRequests::default(),
238 inflight_requests: Mutex::new(InflightRequests::default()),
239 no_inflight_requests_condvar: Condvar::new(),
240 inner: Mutex::new(SessionManagerInner { open_sessions: HashMap::new() }),
241 no_open_sessions_condvar: Condvar::new(),
242 block_service: OnceLock::new(),
243 }
244 }
245
246 pub fn block_size(&self) -> u32 {
247 self.block_size
248 }
249
250 pub fn complete_request(&self, request_id: RequestId, status: Result<(), zx::Status>) {
252 let (callback, notify) = {
253 let mut inflight = self.inflight_requests.lock();
254 let callback = if request_id.0 & INTERNAL_REQUEST_FLAG != 0 {
255 inflight.callbacks.remove(&request_id)
256 } else {
257 None
258 };
259 inflight.count -= 1;
260 let notify = inflight.count == 0;
261 (callback, notify)
262 };
263 if let Some(callback) = callback {
264 callback(status);
265 } else {
266 self.complete_unsubmitted_request(request_id, status);
267 }
268 if notify {
269 self.no_inflight_requests_condvar.notify_all();
270 }
271 }
272
273 fn submit_internal_request(
274 &self,
275 make_request: impl FnOnce(RequestId) -> Request,
276 callback: Box<dyn FnOnce(Result<(), zx::Status>) + Send>,
277 ) {
278 let req = {
279 let mut inflight = self.inflight_requests.lock();
280 let req_id = RequestId(inflight.next_internal_id | INTERNAL_REQUEST_FLAG);
281 inflight.next_internal_id += 1;
282 inflight.callbacks.insert(req_id, callback);
283 inflight.count += 1;
284 make_request(req_id)
285 };
286 self.interface.on_requests(&[req]);
287 }
288
289 fn submit_requests(&self, requests: &[Request]) {
290 self.inflight_requests.lock().count += requests.len();
291 self.interface.on_requests(requests);
292 }
293
294 fn wait_for_no_inflight_requests(&self) {
299 let mut guard = self.inflight_requests.lock();
300 self.no_inflight_requests_condvar.wait_while(&mut guard, |inflight| inflight.count > 0);
301 }
302
303 fn complete_unsubmitted_request(&self, request_id: RequestId, status: Result<(), zx::Status>) {
306 if let Some((session, response)) =
307 self.active_requests.complete_and_take_response(request_id, status)
308 {
309 session.send_response(response);
310 }
311 }
312
313 pub fn terminate(&self) {
315 {
316 #[allow(clippy::collection_is_never_read)]
319 let mut terminated_sessions = Vec::new();
320 for (_, session) in &self.inner.lock().open_sessions {
321 if let Some(session) = session.upgrade() {
322 session.terminate_async();
323 terminated_sessions.push(session);
324 }
325 }
326 }
327 let mut guard = self.inner.lock();
328 self.no_open_sessions_condvar.wait_while(&mut guard, |s| !s.open_sessions.is_empty());
329 }
330}
331
332pub struct Session<I: Interface + ?Sized> {
333 helper: SessionHelper<SessionManager<I>>,
334 fifo: zx::Fifo<BlockFifoRequest, BlockFifoResponse>,
335 queue: Mutex<SessionQueue>,
336 abort_handle: AbortHandle,
337 close_callback: Mutex<Option<Box<dyn FnOnce() + Send + 'static>>>,
338}
339
340#[derive(Default)]
341struct SessionQueue {
342 responses: VecDeque<BlockFifoResponse>,
343}
344
345pub const MAX_REQUESTS: usize = super::FIFO_MAX_REQUESTS;
346
347struct DecodedRequests {
348 requests: [MaybeUninit<Request>; MAX_REQUESTS],
349 count: usize,
350}
351
352impl Default for DecodedRequests {
353 fn default() -> Self {
354 Self { requests: [const { MaybeUninit::uninit() }; MAX_REQUESTS], count: 0 }
355 }
356}
357
358impl DecodedRequests {
359 fn push(&mut self, request: Request) {
360 assert!(self.count < MAX_REQUESTS);
361 self.requests[self.count].write(request);
364 self.count += 1;
365 }
366
367 fn is_empty(&self) -> bool {
368 self.count == 0
369 }
370
371 fn is_full(&self) -> bool {
372 self.count == MAX_REQUESTS
373 }
374
375 fn clear(&mut self) {
376 for i in 0..self.count {
377 unsafe { self.requests[i].assume_init_drop() };
379 }
380 self.count = 0;
381 }
382}
383
384impl Drop for DecodedRequests {
385 fn drop(&mut self) {
386 self.clear();
387 }
388}
389
390impl std::ops::Deref for DecodedRequests {
391 type Target = [Request];
392
393 fn deref(&self) -> &Self::Target {
394 unsafe { std::slice::from_raw_parts(self.requests[0].as_ptr(), self.count) }
396 }
397}
398
399impl std::ops::DerefMut for DecodedRequests {
400 fn deref_mut(&mut self) -> &mut Self::Target {
401 unsafe { std::slice::from_raw_parts_mut(self.requests[0].as_mut_ptr(), self.count) }
403 }
404}
405
406impl<I: Interface + ?Sized> Session<I> {
407 pub fn run(self: &Arc<Self>) {
410 self.fifo_loop();
411 self.abort_handle.abort();
412 self.helper.close_active_groups(|s| Arc::ptr_eq(s, self));
416 }
417
418 fn fifo_loop(self: &Arc<Self>) {
419 let mut requests = [MaybeUninit::uninit(); MAX_REQUESTS];
420
421 loop {
422 let is_queue_empty = {
424 let mut queue = self.queue.lock();
425 while !queue.responses.is_empty() {
426 let (front, _) = queue.responses.as_slices();
427 match self.fifo.write(front) {
428 Ok(count) => {
429 let full = count.get() < front.len();
430 queue.responses.drain(..count.get());
431 if full {
432 break;
433 }
434 }
435 Err(zx::Status::SHOULD_WAIT) => break,
436 Err(_) => return,
437 }
438 }
439 queue.responses.is_empty()
440 };
441
442 match self.fifo.read_uninit(&mut requests) {
444 Ok(valid_requests) => self.handle_requests(valid_requests.iter_mut()),
445 Err(zx::Status::SHOULD_WAIT) => {
446 let mut signals =
447 zx::Signals::OBJECT_READABLE | SHUTDOWN_SIGNAL | FIFO_WAKE_SIGNAL;
448 if !is_queue_empty {
449 signals |= zx::Signals::OBJECT_WRITABLE;
450 }
451 let Ok(signals) =
452 self.fifo.wait_one(signals, zx::MonotonicInstant::INFINITE).to_result()
453 else {
454 return;
455 };
456 if signals.contains(SHUTDOWN_SIGNAL) {
457 return;
458 }
459 if signals.contains(FIFO_WAKE_SIGNAL) {
461 let _ = self.fifo.signal(FIFO_WAKE_SIGNAL, zx::Signals::empty());
462 }
463 }
464 Err(_) => return,
465 }
466 }
467 }
468
469 fn pre_flush(self: &Arc<Self>, request_id: RequestId) -> Result<(), zx::Status> {
471 let trace_flow_id = {
472 let mut request = self.helper.session_manager().active_requests.request(request_id);
473 if let Some(id) = request.trace_flow_id {
474 fuchsia_trace::async_instant!(
475 fuchsia_trace::Id::from(id.get()),
476 c"storage",
477 c"block_server::SimulatedBarrier",
478 "request_id" => request_id.0
479 );
480 }
481 request.count += 1;
482 request.trace_flow_id
483 };
484 self.helper.session_manager().submit_requests(&[Request {
485 request_id,
486 operation: Operation::Flush,
487 trace_flow_id,
488 vmo: None,
489 }]);
490 self.helper.session_manager().wait_for_no_inflight_requests();
491 let status = self.helper.session_manager().active_requests.request(request_id).status;
492 match status {
493 Ok(()) => Ok(()),
494 Err(status) => {
495 self.helper.session_manager().complete_unsubmitted_request(request_id, Err(status));
497 Err(status)
498 }
499 }
500 }
501
502 fn post_flush(self: &Arc<Self>, request_id: RequestId, decoded_requests: &mut DecodedRequests) {
505 if !decoded_requests.is_empty() {
506 self.helper.session_manager().submit_requests(decoded_requests);
507 decoded_requests.clear();
508 }
509 self.helper.session_manager().wait_for_no_inflight_requests();
510 let request = self.helper.session_manager().active_requests.request(request_id);
511 match request.status {
512 Ok(()) => decoded_requests.push(Request {
513 request_id,
514 operation: Operation::Flush,
515 trace_flow_id: request.trace_flow_id,
516 vmo: None,
517 }),
518 Err(status) => {
519 drop(request);
520 self.helper.session_manager().complete_unsubmitted_request(request_id, Err(status))
521 }
522 }
523 }
524
525 fn handle_requests<'a>(
526 self: &Arc<Self>,
527 requests: impl Iterator<Item = &'a mut BlockFifoRequest>,
528 ) {
529 let manager = &self.helper.session_manager();
530 let mut decoded_requests = DecodedRequests::default();
531
532 for request in requests {
533 match self.helper.decode_fifo_request(self.clone(), request) {
534 Ok(DecodedRequest { operation: Operation::CloseVmo, request_id, .. }) => {
535 manager.complete_unsubmitted_request(request_id, Ok(()));
536 }
537 Ok(mut request) => {
538 let request_id = request.request_id;
539
540 if !manager
543 .interface
544 .get_info()
545 .device_flags()
546 .contains(fblock::DeviceFlag::BARRIER_SUPPORT)
547 && request.operation.take_write_flag(WriteFlags::PRE_BARRIER)
548 && self.pre_flush(request_id).is_err()
549 {
550 continue;
551 }
552 let simulate_fua = !manager
555 .interface
556 .get_info()
557 .device_flags()
558 .contains(fblock::DeviceFlag::FUA_SUPPORT)
559 && request.operation.take_write_flag(WriteFlags::FORCE_ACCESS);
560
561 if simulate_fua {
562 manager.active_requests.request(request_id).count += 1;
564 }
565
566 loop {
567 let result = self
568 .helper
569 .map_request(request, &mut manager.active_requests.request(request_id));
570 match result {
571 Ok((
572 DecodedRequest { request_id, operation, vmo, trace_flow_id },
573 remainder,
574 )) => {
575 decoded_requests.push(Request {
576 request_id,
577 operation,
578 trace_flow_id,
579 vmo,
580 });
581
582 if decoded_requests.is_full() {
583 manager.submit_requests(&*decoded_requests);
584 decoded_requests.clear();
585 }
586
587 if let Some(r) = remainder {
588 request = r;
589 } else {
590 break;
591 }
592 }
593 Err(status) => {
594 manager.complete_unsubmitted_request(request_id, Err(status));
595 break;
596 }
597 }
598 }
599
600 if simulate_fua {
601 self.post_flush(request_id, &mut decoded_requests);
602 }
603 }
604 Err(None) => {}
605 Err(Some(response)) => self.send_response(response),
606 }
607 }
608
609 if !decoded_requests.is_empty() {
610 manager.submit_requests(&decoded_requests);
611 }
612 }
613
614 fn send_response(&self, response: BlockFifoResponse) {
615 let mut queue = self.queue.lock();
616 if queue.responses.is_empty() {
617 match self.fifo.write_one(&response) {
618 Ok(()) => {
619 return;
620 }
621 Err(_) => {
622 let _ = self.fifo.signal(zx::Signals::empty(), FIFO_WAKE_SIGNAL);
624 }
625 }
626 }
627 queue.responses.push_back(response);
628 }
629
630 pub fn terminate_async(&self) {
633 let _ = self.fifo.signal(zx::Signals::empty(), SHUTDOWN_SIGNAL);
634 self.abort_handle.abort();
635 }
636}
637
638impl<I: Interface + ?Sized> Drop for Session<I> {
639 fn drop(&mut self) {
640 let callback = std::mem::take(&mut *self.close_callback.lock());
641 if let Some(callback) = callback {
642 callback();
643 }
644 let notify = {
645 let mut inner = self.helper.session_manager().inner.lock();
646 inner.open_sessions.remove(&(self as *const _ as usize));
647 inner.open_sessions.is_empty()
648 };
649 if notify {
650 self.helper.session_manager().no_open_sessions_condvar.notify_all();
651 }
652 }
653}
654
655impl<I: Interface + ?Sized> Drop for SessionManager<I> {
656 fn drop(&mut self) {
657 self.terminate();
658 }
659}
660
661impl<I: Interface<Orchestrator = SessionManager<I>>> IntoOrchestrator for Arc<SessionManager<I>> {
662 type SM = SessionManager<I>;
663
664 fn into_orchestrator(self) -> Arc<I::Orchestrator> {
665 self
666 }
667}
668
669#[cfg(test)]
670mod tests {
671 use super::*;
672 use crate::testing::MockInterface;
673 use block_protocol::{BlockFifoCommand, BlockFifoRequest, BlockFifoResponse};
674 use fidl::endpoints::create_proxy_and_stream;
675 use fidl_fuchsia_storage_block as fblock;
676 use fuchsia_async as fasync;
677
678 const BLOCK_SIZE: u32 = 512;
679
680 #[fuchsia::test]
681 async fn test_basic_request() {
682 let (tx, rx) = std::sync::mpsc::channel();
683 let interface = Arc::new(MockInterface::new(tx));
684 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
685
686 let sm_clone = session_manager.clone();
687 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
688 let _server_task = fasync::Task::spawn(async move {
689 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
690 server.handle_requests(stream).await.unwrap();
691 })
692 .detach();
693
694 let (session_proxy, session_server_end) =
695 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
696 proxy.open_session(session_server_end).unwrap();
697
698 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
699 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
700
701 let vmo = zx::Vmo::create(8192).unwrap();
702 let vmo_id = session_proxy
703 .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
704 .await
705 .unwrap()
706 .unwrap();
707
708 let req = BlockFifoRequest {
709 command: BlockFifoCommand {
710 opcode: fblock::BlockOpcode::Read.into_primitive(),
711 ..Default::default()
712 },
713 reqid: 123,
714 group: 0,
715 vmoid: vmo_id.id,
716 length: 1,
717 vmo_offset: 0,
718 dev_offset: 0,
719 trace_flow_id: 0,
720 ..Default::default()
721 };
722 fifo.write(&[req]).unwrap();
723
724 let r = rx.recv().unwrap();
725 assert_eq!(r.request_id.0, 0);
726 assert!(matches!(r.operation, Operation::Read { .. }));
727
728 session_manager.complete_request(r.request_id, Ok(()));
729
730 let signals =
731 fifo.wait_one(zx::Signals::FIFO_READABLE, zx::MonotonicInstant::INFINITE).unwrap();
732 assert!(signals.contains(zx::Signals::FIFO_READABLE));
733
734 let mut resp = [BlockFifoResponse::default()];
735 fifo.read(&mut resp).unwrap();
736 assert_eq!(resp[0].reqid, 123);
737 assert_eq!(resp[0].status, zx::sys::ZX_OK);
738
739 std::mem::drop(proxy);
740 }
741 #[fuchsia::test]
742 async fn test_write_request() {
743 let (tx, rx) = std::sync::mpsc::channel();
744 let interface = Arc::new(MockInterface::new(tx));
745 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
746
747 let sm_clone = session_manager.clone();
748 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
749 let _server_task = fasync::Task::spawn(async move {
750 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
751 server.handle_requests(stream).await.unwrap();
752 })
753 .detach();
754
755 let (session_proxy, session_server_end) =
756 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
757 proxy.open_session(session_server_end).unwrap();
758
759 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
760 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
761
762 let vmo = zx::Vmo::create(8192).unwrap();
763 let vmo_id = session_proxy
764 .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
765 .await
766 .unwrap()
767 .unwrap();
768
769 let req = BlockFifoRequest {
770 command: BlockFifoCommand {
771 opcode: fblock::BlockOpcode::Write.into_primitive(),
772 ..Default::default()
773 },
774 reqid: 124,
775 group: 0,
776 vmoid: vmo_id.id,
777 length: 1,
778 vmo_offset: 0,
779 dev_offset: 0,
780 trace_flow_id: 0,
781 ..Default::default()
782 };
783 fifo.write(&[req]).unwrap();
784
785 let r = rx.recv().unwrap();
786 assert_eq!(r.request_id.0, 0);
787 assert!(matches!(r.operation, Operation::Write { .. }));
788
789 session_manager.complete_request(r.request_id, Ok(()));
790
791 let signals =
792 fifo.wait_one(zx::Signals::FIFO_READABLE, zx::MonotonicInstant::INFINITE).unwrap();
793 assert!(signals.contains(zx::Signals::FIFO_READABLE));
794
795 let mut resp = [BlockFifoResponse::default()];
796 fifo.read(&mut resp).unwrap();
797 assert_eq!(resp[0].reqid, 124);
798 assert_eq!(resp[0].status, zx::sys::ZX_OK);
799 }
800
801 #[fuchsia::test]
802 async fn test_flush_request() {
803 let (tx, rx) = std::sync::mpsc::channel();
804 let interface = Arc::new(MockInterface::new(tx));
805 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
806
807 let sm_clone = session_manager.clone();
808 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
809 let _server_task = fasync::Task::spawn(async move {
810 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
811 server.handle_requests(stream).await.unwrap();
812 })
813 .detach();
814
815 let (session_proxy, session_server_end) =
816 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
817 proxy.open_session(session_server_end).unwrap();
818
819 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
820 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
821
822 let req = BlockFifoRequest {
823 command: BlockFifoCommand {
824 opcode: fblock::BlockOpcode::Flush.into_primitive(),
825 ..Default::default()
826 },
827 reqid: 125,
828 group: 0,
829 vmoid: fblock::VMOID_INVALID,
830 length: 0,
831 vmo_offset: 0,
832 dev_offset: 0,
833 trace_flow_id: 0,
834 ..Default::default()
835 };
836 fifo.write(&[req]).unwrap();
837
838 let r = rx.recv().unwrap();
839 assert_eq!(r.request_id.0, 0);
840 assert!(matches!(r.operation, Operation::Flush { .. }));
841
842 session_manager.complete_request(r.request_id, Ok(()));
843
844 let signals =
845 fifo.wait_one(zx::Signals::FIFO_READABLE, zx::MonotonicInstant::INFINITE).unwrap();
846 assert!(signals.contains(zx::Signals::FIFO_READABLE));
847
848 let mut resp = [BlockFifoResponse::default()];
849 fifo.read(&mut resp).unwrap();
850 assert_eq!(resp[0].reqid, 125);
851 assert_eq!(resp[0].status, zx::sys::ZX_OK);
852 }
853
854 #[fuchsia::test]
855 async fn test_trim_request() {
856 let (tx, rx) = std::sync::mpsc::channel();
857 let interface = Arc::new(MockInterface::new(tx));
858 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
859
860 let sm_clone = session_manager.clone();
861 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
862 let _server_task = fasync::Task::spawn(async move {
863 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
864 server.handle_requests(stream).await.unwrap();
865 })
866 .detach();
867
868 let (session_proxy, session_server_end) =
869 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
870 proxy.open_session(session_server_end).unwrap();
871
872 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
873 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
874
875 let req = BlockFifoRequest {
876 command: BlockFifoCommand {
877 opcode: fblock::BlockOpcode::Trim.into_primitive(),
878 ..Default::default()
879 },
880 reqid: 126,
881 group: 0,
882 vmoid: fblock::VMOID_INVALID,
883 length: 1,
884 vmo_offset: 0,
885 dev_offset: 0,
886 trace_flow_id: 0,
887 ..Default::default()
888 };
889 fifo.write(&[req]).unwrap();
890
891 let r = rx.recv().unwrap();
892 assert_eq!(r.request_id.0, 0);
893 assert!(matches!(r.operation, Operation::Trim { .. }));
894
895 session_manager.complete_request(r.request_id, Ok(()));
896
897 let signals =
898 fifo.wait_one(zx::Signals::FIFO_READABLE, zx::MonotonicInstant::INFINITE).unwrap();
899 assert!(signals.contains(zx::Signals::FIFO_READABLE));
900
901 let mut resp = [BlockFifoResponse::default()];
902 fifo.read(&mut resp).unwrap();
903 assert_eq!(resp[0].reqid, 126);
904 assert_eq!(resp[0].status, zx::sys::ZX_OK);
905 }
906
907 #[fuchsia::test]
908 async fn test_close_vmo() {
909 let (tx, rx) = std::sync::mpsc::channel();
910 let interface = Arc::new(MockInterface::new(tx));
911 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
912
913 let sm_clone = session_manager.clone();
914 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
915 let _server_task = fasync::Task::spawn(async move {
916 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
917 server.handle_requests(stream).await.unwrap();
918 })
919 .detach();
920
921 let (session_proxy, session_server_end) =
922 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
923 proxy.open_session(session_server_end).unwrap();
924
925 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
926 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
927
928 let vmo = zx::Vmo::create(8192).unwrap();
929 let vmo_id = session_proxy
930 .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
931 .await
932 .unwrap()
933 .unwrap();
934
935 let req = BlockFifoRequest {
936 command: BlockFifoCommand {
937 opcode: fblock::BlockOpcode::CloseVmo.into_primitive(),
938 ..Default::default()
939 },
940 reqid: 127,
941 group: 0,
942 vmoid: vmo_id.id,
943 length: 0,
944 vmo_offset: 0,
945 dev_offset: 0,
946 trace_flow_id: 0,
947 ..Default::default()
948 };
949 fifo.write(&[req]).unwrap();
950
951 let signals =
952 fifo.wait_one(zx::Signals::FIFO_READABLE, zx::MonotonicInstant::INFINITE).unwrap();
953 assert!(signals.contains(zx::Signals::FIFO_READABLE));
954
955 let mut resp = [BlockFifoResponse::default()];
956 fifo.read(&mut resp).unwrap();
957 assert_eq!(resp[0].reqid, 127);
958 assert_eq!(resp[0].status, zx::sys::ZX_OK);
959
960 assert!(rx.try_recv().is_err());
962 }
963
964 #[fuchsia::test]
965 async fn test_error() {
966 let (tx, rx) = std::sync::mpsc::channel();
967 let interface = Arc::new(MockInterface::new(tx));
968 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
969
970 let sm_clone = session_manager.clone();
971 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
972 let _server_task = fasync::Task::spawn(async move {
973 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
974 server.handle_requests(stream).await.unwrap();
975 })
976 .detach();
977
978 let (session_proxy, session_server_end) =
979 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
980 proxy.open_session(session_server_end).unwrap();
981
982 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
983 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
984
985 let req = BlockFifoRequest {
986 command: BlockFifoCommand {
987 opcode: fblock::BlockOpcode::Flush.into_primitive(),
988 ..Default::default()
989 },
990 reqid: 128,
991 group: 0,
992 vmoid: fblock::VMOID_INVALID,
993 length: 0,
994 vmo_offset: 0,
995 dev_offset: 0,
996 trace_flow_id: 0,
997 ..Default::default()
998 };
999 fifo.write(&[req]).unwrap();
1000
1001 let r = rx.recv().unwrap();
1002 session_manager.complete_request(r.request_id, Err(zx::Status::IO));
1003
1004 let signals =
1005 fifo.wait_one(zx::Signals::FIFO_READABLE, zx::MonotonicInstant::INFINITE).unwrap();
1006 assert!(signals.contains(zx::Signals::FIFO_READABLE));
1007
1008 let mut resp = [BlockFifoResponse::default()];
1009 fifo.read(&mut resp).unwrap();
1010 assert_eq!(resp[0].reqid, 128);
1011 assert_eq!(resp[0].status, zx::sys::ZX_ERR_IO);
1012 }
1013
1014 #[fuchsia::test]
1015 async fn test_teardown_with_active_requests() {
1016 let (tx, rx) = std::sync::mpsc::channel();
1017 let interface = Arc::new(MockInterface::new(tx));
1018 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
1019
1020 let sm_clone = session_manager.clone();
1021 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
1022 let server_task = fasync::Task::spawn(async move {
1023 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
1024 server.handle_requests(stream).await.unwrap();
1025 });
1026
1027 let (session_proxy, session_server_end) =
1028 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
1029 proxy.open_session(session_server_end).unwrap();
1030
1031 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
1032 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
1033
1034 let vmo = zx::Vmo::create(8192).unwrap();
1035 let vmo_id = session_proxy
1036 .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
1037 .await
1038 .unwrap()
1039 .unwrap();
1040
1041 let req = BlockFifoRequest {
1042 command: BlockFifoCommand {
1043 opcode: fblock::BlockOpcode::Read.into_primitive(),
1044 ..Default::default()
1045 },
1046 reqid: 129,
1047 group: 0,
1048 vmoid: vmo_id.id,
1049 length: 1,
1050 vmo_offset: 0,
1051 dev_offset: 0,
1052 trace_flow_id: 0,
1053 ..Default::default()
1054 };
1055 fifo.write(&[req]).unwrap();
1057 let r = rx.recv().unwrap();
1058
1059 drop(session_proxy);
1061 fasync::Timer::new(std::time::Duration::from_millis(50)).await;
1062
1063 session_manager.complete_request(r.request_id, Ok(()));
1065
1066 drop(proxy);
1067 fasync::unblock(move || session_manager.terminate()).await;
1068 server_task.await;
1069 }
1070
1071 #[fuchsia::test]
1072 async fn test_teardown_with_active_grouped_requests() {
1073 let (tx, rx) = std::sync::mpsc::channel();
1074 let interface = Arc::new(MockInterface::new(tx));
1075 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
1076
1077 let sm_clone = session_manager.clone();
1078 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
1079 let server_task = fasync::Task::spawn(async move {
1080 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
1081 server.handle_requests(stream).await.unwrap();
1082 });
1083
1084 let (session_proxy, session_server_end) =
1085 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
1086 proxy.open_session(session_server_end).unwrap();
1087
1088 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
1089 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
1090
1091 let vmo = zx::Vmo::create(8192).unwrap();
1092 let vmo_id = session_proxy
1093 .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
1094 .await
1095 .unwrap()
1096 .unwrap();
1097
1098 let mut req = BlockFifoRequest {
1099 command: BlockFifoCommand {
1100 opcode: fblock::BlockOpcode::Read.into_primitive(),
1101 flags: fblock::BlockIoFlag::GROUP_ITEM.bits(),
1102 ..Default::default()
1103 },
1104 reqid: 1,
1105 group: 1,
1106 vmoid: vmo_id.id,
1107 length: 1,
1108 vmo_offset: 0,
1109 dev_offset: 0,
1110 trace_flow_id: 0,
1111 ..Default::default()
1112 };
1113
1114 fifo.write(&[req]).unwrap();
1116 req.reqid = 2;
1117 req.group = 2;
1118 fifo.write(&[req]).unwrap();
1119 let r1 = rx.recv().unwrap();
1120 let r2 = rx.recv().unwrap();
1121 req.reqid = 3;
1122 req.group = 2;
1123 fifo.write(&[req]).unwrap();
1124 let r3 = rx.recv().unwrap();
1125
1126 session_manager.complete_request(r1.request_id, Ok(()));
1131 session_manager.complete_request(r2.request_id, Ok(()));
1132
1133 drop(session_proxy);
1136 fasync::Timer::new(std::time::Duration::from_millis(50)).await;
1137
1138 session_manager.complete_request(r3.request_id, Ok(()));
1140
1141 drop(proxy);
1142
1143 fasync::unblock(move || session_manager.terminate()).await;
1144 server_task.await;
1145 }
1146
1147 #[fuchsia::test]
1148 async fn test_session_close_is_synchronous() {
1149 use futures::FutureExt as _;
1150
1151 let (tx, rx) = std::sync::mpsc::channel();
1152 let interface = Arc::new(MockInterface::new(tx));
1153 let session_manager = Arc::new(SessionManager::new(interface.clone(), BLOCK_SIZE));
1154
1155 let sm_clone = session_manager.clone();
1156 let (proxy, stream) = create_proxy_and_stream::<fblock::BlockMarker>();
1157 let server_task = fasync::Task::spawn(async move {
1158 let server = crate::BlockServer::new(BLOCK_SIZE, sm_clone);
1159 server.handle_requests(stream).await.unwrap();
1160 });
1161
1162 let (session_proxy, session_server_end) =
1163 fidl::endpoints::create_proxy::<fblock::SessionMarker>();
1164 proxy.open_session(session_server_end).unwrap();
1165
1166 let fifo_handle = session_proxy.get_fifo().await.unwrap().unwrap();
1167 let fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest> = zx::Fifo::from(fifo_handle);
1168
1169 let vmo = zx::Vmo::create(8192).unwrap();
1170 let vmo_id = session_proxy
1171 .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
1172 .await
1173 .unwrap()
1174 .unwrap();
1175
1176 let req = BlockFifoRequest {
1177 command: BlockFifoCommand {
1178 opcode: fblock::BlockOpcode::Read.into_primitive(),
1179 ..Default::default()
1180 },
1181 reqid: 123,
1182 group: 0,
1183 vmoid: vmo_id.id,
1184 length: 1,
1185 vmo_offset: 0,
1186 dev_offset: 0,
1187 trace_flow_id: 0,
1188 ..Default::default()
1189 };
1190 fifo.write(&[req]).unwrap();
1191
1192 let r = rx.recv().unwrap();
1193
1194 let mut close_fut = std::pin::pin!(session_proxy.close().fuse());
1196 let mut timer_fut =
1197 std::pin::pin!(fasync::Timer::new(std::time::Duration::from_millis(100)).fuse());
1198 futures::select! {
1199 res = close_fut => panic!("close completed too early: {res:?}"),
1200 _ = timer_fut => {}
1201 }
1202
1203 session_manager.complete_request(r.request_id, Ok(()));
1204
1205 close_fut.await.unwrap().unwrap();
1207
1208 std::mem::drop(proxy);
1209 server_task.await;
1210 }
1211}