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