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