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 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/// An in-flight request.
25#[derive(Clone, Debug)]
26pub struct Request {
27    /// The ID that a request is associated with, for later completion in
28    /// [`SessionManager::complete_request`].  Note that this is not necessarily a unique
29    /// identifier, and multiple Requests may have the same request_id (e.g. due to request
30    /// splitting).  The library internally reference-counts requests which use this ID.
31    pub request_id: RequestId,
32    pub operation: Operation,
33    pub trace_flow_id: TraceFlowId,
34    /// `vmo` is always Some for Operation::Read or Operation::Write, and None otherwise.
35    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    /// Called to get block/partition information.
42    fn get_info(&self) -> Cow<'_, DeviceInfo>;
43
44    /// Start running a new session.  The interface must ensure that [`Session::run`] is called
45    /// on a different thread.
46    fn spawn_session(&self, session: Arc<Session<Self>>);
47
48    /// Starts a batch of requests.  The implementation may block if there are too many in-flight
49    /// requests, providing pushback.  The interface is responsible for eventually calling
50    /// [`SessionManager::complete_request`] for each request (even during server shutdown).
51    ///
52    /// Implementations are responsible for checking that request block ranges fall within valid
53    /// device/partition bounds, and completing with `zx::Status::OUT_OF_RANGE` if out of bounds.
54    fn on_requests(&self, requests: &[Request]);
55
56    /// Returns the BlockService implementation for this interface.
57    ///
58    /// By default, returns the fallback [`DefaultCallbackBlockService`].
59    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
67/// A bit flag set on [`RequestId`] values for internal requests issued by the server (such as
68/// [`DefaultCallbackBlockService`] reads) rather than external client sessions. This allows
69/// [`SessionManager::complete_request`] to distinguish internal requests and route completions
70/// to internal callback handlers instead of client session FIFOs.
71const INTERNAL_REQUEST_FLAG: usize = 1 << 63;
72
73#[derive(Default)]
74struct InflightRequests {
75    /// Total number of in-flight requests submitted to `Interface::on_requests`.
76    count: usize,
77    /// Completion callbacks for internal requests (e.g. from `DefaultCallbackBlockService`).
78    callbacks: HashMap<RequestId, Box<dyn FnOnce(zx::Status) + Send>>,
79    /// Monotonically increasing counter for allocating internal request IDs.
80    next_internal_id: usize,
81}
82
83struct SessionManagerInner<I: Interface + ?Sized> {
84    open_sessions: HashMap<usize, Weak<Session<I>>>,
85}
86
87// The signals used on the session's FIFO.
88
89/// Signalled on the session's FIFO to wake up the FIFO loop.
90const FIFO_WAKE_SIGNAL: zx::Signals = zx::Signals::USER_0;
91
92/// Signalled on the session's FIFO to terminate the FIFO loop.
93const SHUTDOWN_SIGNAL: zx::Signals = zx::Signals::USER_1;
94
95pub struct SessionManager<I: Interface + ?Sized> {
96    interface: Arc<I>,
97    block_size: u32,
98    // These represent active *client* requests, which correspond to one or more in-flight requests.
99    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    /// Called when a new session request stream is opened. Constructs a `SessionHelper` with the
125    /// client-provided `offset_map` and device constraints, spawns the session runner via
126    /// `I::spawn_session`, and processes FIFO requests until the session channel is closed or
127    /// aborted.
128    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    /// Returns the memoized [`BlockService`] for this session manager.
188    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    /// Reports the given task as complete with a given status.
217    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    /// Waits for there to be no requests in-flight.
261    ///
262    /// NOTE: To void TOCTOUs, this must be called on the same thread which calls
263    /// [`Self::submit_requests`].
264    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    /// Called instead of `[Self::complete_request]` when a request is completed before it was
270    /// actually submitted.
271    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    /// Terminates the session manager.  Blocks until all sessions have terminated.
280    pub fn terminate(&self) {
281        {
282            // We must drop references to sessions whilst we're not holding the lock for
283            // `open_sessions` because `Session::drop` needs to take that same lock.
284            #[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        // To ensure drop runs, we must initialize each element at most once.
328        // As long as we only go via this function, that's satisfied.
329        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            // SAFETY: We initialized `count` elements via [`Self::push`].
344            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        // SAFETY: We wrote the request in [`Self::push`].
361        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        // SAFETY: We wrote the request in [`Self::push`].
368        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    /// Begins processing the session's FIFO.  Blocks until completion (either because of
374    /// termination, or due to unrecoverable error).
375    pub fn run(self: &Arc<Self>) {
376        self.fifo_loop();
377        self.abort_handle.abort();
378        // NB: We cannot call [`drop_active_requests`] here, because requests which have already
379        // been submitted are no longer in control by this thread, and will be later completed by
380        // another thread.  If we dropped them here, then they would be completed twice.
381        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            // Send queued responses.
389            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            // Process pending reads.
409            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                    // Clear FIFO_WAKE_SIGNAL if it's set.
426                    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    /// Synchronously performs a device flush.
436    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                // Respond for the unsubmitted request too.
462                self.helper.session_manager().complete_unsubmitted_request(request_id, status);
463                Err(status)
464            }
465        }
466    }
467
468    /// Synchronously completes `decoded_requests`, and inserts a post-flush into `decoded_requests`
469    /// to be submitted later.
470    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                    // Strip the PRE_BARRIER flag if we don't support it, and simulate the barrier
507                    // with a pre-flush.
508                    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                    // Strip the FORCE_ACCESS flag if we don't support it, and simulate the FUA with
519                    // a post-flush.
520                    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                        // Account for the additional request we need at the end.
529                        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                    // Wake the FIFO loop up so that we can send the response later.
589                    let _ = self.fifo.signal(zx::Signals::empty(), FIFO_WAKE_SIGNAL);
590                }
591            }
592        }
593        queue.responses.push_back(response);
594    }
595
596    /// Asynchronously request to terminate the Session.  The session's thread will eventually stop
597    /// running.
598    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        // CloseVmo is handled automatically. Interface should not receive anything.
951        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        // Start a request and wait until it's been submitted.
1046        fifo.write(&[req]).unwrap();
1047        let r = rx.recv().unwrap();
1048
1049        // Close the client, which will eventually cause the FIFO loop to exit.
1050        drop(session_proxy);
1051        fasync::Timer::new(std::time::Duration::from_millis(50)).await;
1052
1053        // Complete the request, simulating a completion after the FIFO loop has exited.
1054        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        // Start two groups of requests and wait until they're been submitted.
1105        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        // Complete requests 1,2.  This leaves two groups in the following states:
1117        // - Group 1: 0 active requests, still waiting for END
1118        // - Group 2: 1 active request, still waiting for END
1119        // Neither group will be able to complete yet.
1120        session_manager.complete_request(r1.request_id, zx::Status::OK);
1121        session_manager.complete_request(r2.request_id, zx::Status::OK);
1122
1123        // Close the client, which will eventually cause the FIFO loop to exit.
1124        // Group 1 should complete now.  Group 2 can't yet.
1125        drop(session_proxy);
1126        fasync::Timer::new(std::time::Duration::from_millis(50)).await;
1127
1128        // At some later time, complete request 3, which should complete group 2.
1129        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        // The close request shouldn't complete yet because the read is still hanging.
1185        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        // Verify that close() now completes.
1196        close_fut.await.unwrap().unwrap();
1197
1198        std::mem::drop(proxy);
1199        server_task.await;
1200    }
1201}