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