Skip to main content

block_server/
callback_interface.rs

1// Copyright 2026 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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/// An in-flight request.
21#[derive(Clone, Debug)]
22pub struct Request {
23    /// The ID that a request is associated with, for later completion in
24    /// [`SessionManager::complete_request`].  Note that this is not necessarily a unique
25    /// identifier, and multiple Requests may have the same request_id (e.g. due to request
26    /// splitting).  The library internally reference-counts requests which use this ID.
27    pub request_id: RequestId,
28    pub operation: Operation,
29    pub trace_flow_id: TraceFlowId,
30    /// `vmo` is always Some for Operation::Read or Operation::Write, and None otherwise.
31    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    /// Called to get block/partition information.
38    fn get_info(&self) -> Cow<'_, DeviceInfo>;
39
40    /// Start running a new session.  The interface must ensure that [`Session::run`] is called
41    /// on a different thread.
42    fn spawn_session(&self, session: Arc<Session<Self>>);
43
44    /// Starts a batch of requests.  The implementation may block if there are too many in-flight
45    /// requests, providing pushback.  The interface is responsible for eventually calling
46    /// [`SessionManager::complete_request`] for each request (even during server shutdown).
47    ///
48    /// Implementations are responsible for checking that request block ranges fall within valid
49    /// device/partition bounds, and completing with `zx::Status::OUT_OF_RANGE` if out of bounds.
50    fn on_requests(&self, requests: &[Request]);
51}
52
53struct SessionManagerInner<I: Interface + ?Sized> {
54    open_sessions: HashMap<usize, Weak<Session<I>>>,
55}
56
57// The signals used on the session's FIFO.
58
59/// Signalled on the session's FIFO to wake up the FIFO loop.
60const FIFO_WAKE_SIGNAL: zx::Signals = zx::Signals::USER_0;
61
62/// Signalled on the session's FIFO to terminate the FIFO loop.
63const SHUTDOWN_SIGNAL: zx::Signals = zx::Signals::USER_1;
64
65pub struct SessionManager<I: Interface + ?Sized> {
66    interface: Arc<I>,
67    // These represent active *client* requests, which correspond to one or more in-flight requests.
68    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    /// Called when a new session request stream is opened. Constructs a `SessionHelper` with the
93    /// client-provided `offset_map` and device constraints, spawns the session runner via
94    /// `I::spawn_session`, and processes FIFO requests until the session channel is closed or
95    /// aborted.
96    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    /// Reports the given task as complete with a given status.
167    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    /// Waits for there to be no requests in-flight.
185    ///
186    /// NOTE: To void TOCTOUs, this must be called on the same thread which calls
187    /// [`Self::submit_requests`].
188    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    /// Called instead of `[Self::complete_request]` when a request is completed before it was
194    /// actually submitted.
195    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    /// Terminates the session manager.  Blocks until all sessions have terminated.
204    pub fn terminate(&self) {
205        {
206            // We must drop references to sessions whilst we're not holding the lock for
207            // `open_sessions` because `Session::drop` needs to take that same lock.
208            #[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        // To ensure drop runs, we must initialize each element at most once.
252        // As long as we only go via this function, that's satisfied.
253        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            // SAFETY: We initialized `count` elements via [`Self::push`].
268            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        // SAFETY: We wrote the request in [`Self::push`].
285        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        // SAFETY: We wrote the request in [`Self::push`].
292        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    /// Begins processing the session's FIFO.  Blocks until completion (either because of
298    /// termination, or due to unrecoverable error).
299    pub fn run(self: &Arc<Self>) {
300        self.fifo_loop();
301        self.abort_handle.abort();
302        // NB: We cannot call [`drop_active_requests`] here, because requests which have already
303        // been submitted are no longer in control by this thread, and will be later completed by
304        // another thread.  If we dropped them here, then they would be completed twice.
305        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            // Send queued responses.
313            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            // Process pending reads.
333            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                    // Clear FIFO_WAKE_SIGNAL if it's set.
350                    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    /// Synchronously performs a device flush.
360    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                // Respond for the unsubmitted request too.
386                self.helper.session_manager().complete_unsubmitted_request(request_id, status);
387                Err(status)
388            }
389        }
390    }
391
392    /// Synchronously completes `decoded_requests`, and inserts a post-flush into `decoded_requests`
393    /// to be submitted later.
394    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                    // Strip the PRE_BARRIER flag if we don't support it, and simulate the barrier
431                    // with a pre-flush.
432                    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                    // Strip the FORCE_ACCESS flag if we don't support it, and simulate the FUA with
443                    // a post-flush.
444                    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                        // Account for the additional request we need at the end.
453                        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                    // Wake the FIFO loop up so that we can send the response later.
513                    let _ = self.fifo.signal(zx::Signals::empty(), FIFO_WAKE_SIGNAL);
514                }
515            }
516        }
517        queue.responses.push_back(response);
518    }
519
520    /// Asynchronously request to terminate the Session.  The session's thread will eventually stop
521    /// running.
522    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        // CloseVmo is handled automatically. Interface should not receive anything.
873        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        // Start a request and wait until it's been submitted.
968        fifo.write(&[req]).unwrap();
969        let r = rx.recv().unwrap();
970
971        // Close the client, which will eventually cause the FIFO loop to exit.
972        drop(session_proxy);
973        fasync::Timer::new(std::time::Duration::from_millis(50)).await;
974
975        // Complete the request, simulating a completion after the FIFO loop has exited.
976        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        // Start two groups of requests and wait until they're been submitted.
1027        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        // Complete requests 1,2.  This leaves two groups in the following states:
1039        // - Group 1: 0 active requests, still waiting for END
1040        // - Group 2: 1 active request, still waiting for END
1041        // Neither group will be able to complete yet.
1042        session_manager.complete_request(r1.request_id, zx::Status::OK);
1043        session_manager.complete_request(r2.request_id, zx::Status::OK);
1044
1045        // Close the client, which will eventually cause the FIFO loop to exit.
1046        // Group 1 should complete now.  Group 2 can't yet.
1047        drop(session_proxy);
1048        fasync::Timer::new(std::time::Duration::from_millis(50)).await;
1049
1050        // At some later time, complete request 3, which should complete group 2.
1051        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        // The close request shouldn't complete yet because the read is still hanging.
1107        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        // Verify that close() now completes.
1118        close_fut.await.unwrap().unwrap();
1119
1120        std::mem::drop(proxy);
1121        server_task.await;
1122    }
1123}