Skip to main content

block_server/
lib.rs

1// Copyright 2025 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.
4use anyhow::Error;
5use block_protocol::{BlockFifoRequest, BlockFifoResponse};
6use fblock::{BlockIoFlag, BlockOpcode, MAX_TRANSFER_UNBOUNDED};
7use fidl_fuchsia_storage_block as fblock;
8use fuchsia_async as fasync;
9use fuchsia_async::epoch::{Epoch, EpochGuard};
10use fuchsia_sync::{MappedMutexGuard, Mutex, MutexGuard};
11use futures::{Future, FutureExt as _, TryStreamExt as _};
12use slab::Slab;
13use std::borrow::{Borrow, Cow};
14use std::collections::BTreeMap;
15use std::num::NonZero;
16use std::ops::Range;
17use std::sync::Arc;
18use std::sync::atomic::AtomicU64;
19use storage_device::buffer::Buffer;
20
21pub mod async_interface;
22pub mod c_interface;
23pub mod callback_interface;
24mod mapper;
25pub mod verifier;
26
27#[cfg(test)]
28pub mod testing;
29
30#[cfg(test)]
31mod decompression_tests;
32
33pub(crate) const FIFO_MAX_REQUESTS: usize = 64;
34
35type TraceFlowId = Option<NonZero<u64>>;
36
37#[derive(Clone, Debug)]
38pub enum DeviceInfo {
39    /// A raw non-partition block device.
40    Block(BlockInfo),
41    /// A static partition with fixed physical mappings.
42    Partition(PartitionInfo),
43    /// A dynamic volume whose slice/block count is queried via get_volume_info.
44    Volume(VolumeInfo),
45}
46
47impl DeviceInfo {
48    pub fn label(&self) -> &str {
49        match self {
50            Self::Block(BlockInfo { .. }) => "",
51            Self::Partition(PartitionInfo { name, .. }) => name,
52            Self::Volume(VolumeInfo { name, .. }) => name,
53        }
54    }
55    pub fn device_flags(&self) -> fblock::DeviceFlag {
56        match self {
57            Self::Block(BlockInfo { device_flags, .. }) => *device_flags,
58            Self::Partition(PartitionInfo { device_flags, .. }) => *device_flags,
59            Self::Volume(VolumeInfo { device_flags, .. }) => *device_flags,
60        }
61    }
62
63    /// Returns the block count of the device or partition.
64    /// Returns None for dynamic volumes (whose size needs to be queried from the volume manager).
65    pub fn block_count(&self) -> Option<u64> {
66        match self {
67            Self::Block(BlockInfo { block_count, .. }) => Some(*block_count),
68            Self::Partition(PartitionInfo { block_count, .. }) => Some(*block_count),
69            Self::Volume(VolumeInfo { .. }) => None,
70        }
71    }
72
73    pub fn max_transfer_blocks(&self) -> Option<NonZero<u32>> {
74        match self {
75            Self::Block(BlockInfo { max_transfer_blocks, .. }) => max_transfer_blocks.clone(),
76            Self::Partition(PartitionInfo { max_transfer_blocks, .. }) => {
77                max_transfer_blocks.clone()
78            }
79            Self::Volume(VolumeInfo { max_transfer_blocks, .. }) => max_transfer_blocks.clone(),
80        }
81    }
82
83    fn max_transfer_size(&self, block_size: u32) -> u32 {
84        if let Some(max_blocks) = self.max_transfer_blocks() {
85            max_blocks.get() * block_size
86        } else {
87            MAX_TRANSFER_UNBOUNDED
88        }
89    }
90
91    pub fn type_guid(&self) -> Option<[u8; 16]> {
92        match self {
93            Self::Partition(PartitionInfo { type_guid, .. }) => Some(*type_guid),
94            Self::Volume(VolumeInfo { type_guid, .. }) => Some(*type_guid),
95            Self::Block(_) => None,
96        }
97    }
98
99    pub fn instance_guid(&self) -> Option<[u8; 16]> {
100        match self {
101            Self::Partition(PartitionInfo { instance_guid, .. }) => Some(*instance_guid),
102            Self::Volume(VolumeInfo { instance_guid, .. }) => Some(*instance_guid),
103            Self::Block(_) => None,
104        }
105    }
106}
107
108/// Information associated with non-partition block devices.
109#[derive(Clone, Default, Debug)]
110pub struct BlockInfo {
111    pub device_flags: fblock::DeviceFlag,
112    pub block_count: u64,
113    pub max_transfer_blocks: Option<NonZero<u32>>,
114}
115
116/// Information associated with a block device that is also a partition.
117#[derive(Clone, Default, Debug)]
118pub struct PartitionInfo {
119    /// The device flags reported by the underlying device.
120    pub device_flags: fblock::DeviceFlag,
121    pub max_transfer_blocks: Option<NonZero<u32>>,
122    /// This is None for partitions which have multiple logical extents, in which case the
123    /// start_block_offset is not meaningful and potentially confusing.
124    pub start_block_offset: Option<u64>,
125    pub block_count: u64,
126    pub type_guid: [u8; 16],
127    pub instance_guid: [u8; 16],
128    pub name: String,
129    /// This can be None for partitions which are composed of multiple partitions (e.g. an
130    /// overlay partition in GPT).
131    pub flags: Option<u64>,
132}
133
134/// Information associated with a dynamic volume (such as FVM volumes).
135#[derive(Clone, Default, Debug)]
136pub struct VolumeInfo {
137    /// The device flags reported by the underlying device.
138    pub device_flags: fblock::DeviceFlag,
139    pub max_transfer_blocks: Option<NonZero<u32>>,
140    pub type_guid: [u8; 16],
141    pub instance_guid: [u8; 16],
142    pub name: String,
143    pub flags: u64,
144}
145
146/// We internally keep track of active requests, so that when the server is torn down, we can
147/// deallocate all of the resources for pending requests.
148struct ActiveRequest<S> {
149    session: S,
150    group_or_request: GroupOrRequest,
151    trace_flow_id: TraceFlowId,
152    _epoch_guard: EpochGuard<'static>,
153    status: Result<(), zx::Status>,
154    count: u32,
155    req_id: Option<u32>,
156    decompression_info: Option<DecompressionInfo>,
157}
158
159struct DecompressionInfo {
160    // This is the range of compressed bytes in receiving buffer.
161    compressed_range: Range<usize>,
162
163    // This is the range in the target VMO where we will write uncompressed bytes.
164    uncompressed_range: Range<u64>,
165
166    bytes_so_far: u64,
167    mapping: Arc<VmoMapping>,
168    buffer: Option<Buffer<'static>>,
169}
170
171impl DecompressionInfo {
172    /// Returns the uncompressed slice.
173    fn uncompressed_slice(&self) -> *mut [u8] {
174        std::ptr::slice_from_raw_parts_mut(
175            (self.mapping.base + self.uncompressed_range.start as usize) as *mut u8,
176            (self.uncompressed_range.end - self.uncompressed_range.start) as usize,
177        )
178    }
179}
180
181pub struct ActiveRequests<S>(Mutex<ActiveRequestsInner<S>>);
182
183impl<S> Default for ActiveRequests<S> {
184    fn default() -> Self {
185        Self(Mutex::new(ActiveRequestsInner { requests: Slab::default() }))
186    }
187}
188
189impl<S> ActiveRequests<S> {
190    fn complete_and_take_response(
191        &self,
192        request_id: RequestId,
193        status: Result<(), zx::Status>,
194    ) -> Option<(S, BlockFifoResponse)> {
195        self.0.lock().complete_and_take_response(request_id, status)
196    }
197
198    fn request(&self, request_id: RequestId) -> MappedMutexGuard<'_, ActiveRequest<S>> {
199        MutexGuard::map(self.0.lock(), |i| &mut i.requests[request_id.0])
200    }
201}
202
203struct ActiveRequestsInner<S> {
204    requests: Slab<ActiveRequest<S>>,
205}
206
207// Keeps track of all the requests that are currently being processed
208impl<S> ActiveRequestsInner<S> {
209    /// Completes a request.
210    fn complete(&mut self, request_id: RequestId, status: Result<(), zx::Status>) {
211        let group = &mut self.requests[request_id.0];
212
213        group.count = group.count.checked_sub(1).unwrap();
214        if status.is_err() && group.status.is_ok() {
215            group.status = status
216        }
217
218        fuchsia_trace::duration!(
219            "storage",
220            "block_server::finish_transaction",
221            "request_id" => request_id.0,
222            "group_completed" => group.count == 0,
223            "status" => zx::Status::result_into_raw(status));
224        if let Some(trace_flow_id) = group.trace_flow_id {
225            fuchsia_trace::flow_step!(
226                "storage",
227                "block_server::finish_request",
228                trace_flow_id.get().into()
229            );
230        }
231
232        if group.count == 0
233            && group.status.is_ok()
234            && let Some(info) = &mut group.decompression_info
235        {
236            struct RawDCtx(std::ptr::NonNull<zstd::zstd_safe::zstd_sys::ZSTD_DCtx>);
237
238            thread_local! {
239                static RAW_DECOMPRESSOR: std::cell::RefCell<RawDCtx> = {
240                    // SAFETY: Creating a new ZSTD decompression context does not borrow or capture
241                    // any external memory.
242                    let raw_ptr = unsafe { zstd::zstd_safe::zstd_sys::ZSTD_createDCtx() };
243                    let ptr = std::ptr::NonNull::new(raw_ptr).expect("ZSTD_createDCtx failed");
244                    std::cell::RefCell::new(RawDCtx(ptr))
245                };
246            }
247
248            impl Drop for RawDCtx {
249                fn drop(&mut self) {
250                    // SAFETY: `self.0` is non-null and was allocated via `ZSTD_createDCtx`.
251                    unsafe {
252                        zstd::zstd_safe::zstd_sys::ZSTD_freeDCtx(self.0.as_ptr());
253                    }
254                }
255            }
256
257            RAW_DECOMPRESSOR.with_borrow_mut(|decompressor| {
258                let dctx = decompressor.0.as_ptr();
259                let target = info.uncompressed_slice();
260                let buffer = info.buffer.take().unwrap();
261                let source = buffer.subslice(info.compressed_range.clone());
262
263                // SAFETY: `target` points to valid uncompressed destination memory in the VMO
264                // mapping. `source` points to valid compressed memory within `buffer`.
265                unsafe {
266                    let result = zstd::zstd_safe::zstd_sys::ZSTD_decompressDCtx(
267                        dctx,
268                        target as *mut u8 as *mut std::os::raw::c_void,
269                        target.len(),
270                        source.as_ptr() as *const std::os::raw::c_void,
271                        source.len(),
272                    );
273                    if zstd::zstd_safe::zstd_sys::ZSTD_isError(result) != 0 {
274                        let error = zstd::zstd_safe::get_error_name(result);
275                        log::warn!(error:?; "Decompression error");
276                        group.status = Err(zx::Status::IO_DATA_INTEGRITY);
277                    }
278                }
279            });
280        }
281    }
282
283    /// Takes the response if all requests are finished.
284    fn take_response(&mut self, request_id: RequestId) -> Option<(S, BlockFifoResponse)> {
285        let group = &self.requests[request_id.0];
286        match group.req_id {
287            Some(reqid) if group.count == 0 => {
288                let group = self.requests.remove(request_id.0);
289                Some((
290                    group.session,
291                    BlockFifoResponse {
292                        status: zx::Status::result_into_raw(group.status),
293                        reqid,
294                        group: group.group_or_request.group_id().unwrap_or(0),
295                        ..Default::default()
296                    },
297                ))
298            }
299            _ => None,
300        }
301    }
302
303    /// Competes the request and returns a response if the request group is finished.
304    fn complete_and_take_response(
305        &mut self,
306        request_id: RequestId,
307        status: Result<(), zx::Status>,
308    ) -> Option<(S, BlockFifoResponse)> {
309        self.complete(request_id, status);
310        self.take_response(request_id)
311    }
312}
313
314/// BlockServer is an implementation of fuchsia.hardware.block.partition.Partition.
315/// cbindgen:no-export
316pub struct BlockServer<SM: SessionManager> {
317    block_size: u32,
318    orchestrator: Arc<SM::Orchestrator>,
319}
320
321#[derive(Clone, Debug, PartialEq, Eq)]
322pub struct BlockOffsetMapping {
323    pub target_block_offset: u64,
324    pub length: u64,
325}
326
327/// Merges physically contiguous mappings into single logical mappings.  This is useful because it
328/// prevents requests from being unnecessarily split if they span the two mappings.
329pub fn coalesce_mappings(raw_mappings: Vec<BlockOffsetMapping>) -> Vec<BlockOffsetMapping> {
330    let mut mappings: Vec<BlockOffsetMapping> = Vec::with_capacity(raw_mappings.len());
331    for m in raw_mappings {
332        if let Some(last) = mappings.last_mut()
333            && last.target_block_offset + last.length == m.target_block_offset
334        {
335            last.length += m.length;
336        } else {
337            mappings.push(m);
338        }
339    }
340    mappings
341}
342
343/// Remaps the offset of block requests based on an internal map of contiguous logical extents.
344#[derive(Clone, Debug, Default, PartialEq, Eq)]
345pub struct OffsetMap {
346    mappings: Vec<BlockOffsetMapping>,
347}
348
349impl OffsetMap {
350    /// Creates a new `OffsetMap` from a list of `BlockOffsetMapping`s.
351    /// Returns `INVALID_ARGS` if any mapping length is zero, or if logical/target block offsets
352    /// overflow `u64`.
353    pub fn new(mappings: Vec<BlockOffsetMapping>) -> Result<Self, zx::Status> {
354        let mut total: u64 = 0;
355        for m in &mappings {
356            if m.length == 0 {
357                return Err(zx::Status::INVALID_ARGS);
358            }
359            m.target_block_offset.checked_add(m.length).ok_or(zx::Status::INVALID_ARGS)?;
360            total = total.checked_add(m.length).ok_or(zx::Status::INVALID_ARGS)?;
361        }
362        Ok(Self { mappings })
363    }
364
365    /// Creates an empty `OffsetMap`.
366    pub fn empty() -> Self {
367        Self { mappings: Vec::new() }
368    }
369
370    /// Maps `logical_offset` to `(target_block_offset, len)`, where `len` is the total number of
371    /// blocks which can be addressed from `target_block_offset` onwards.
372    ///
373    /// For example: if you have logical offset 0 pointing to physical extent 1000..1010, and you
374    /// search for offset 5, the return value will be `Some((1005, 5))`.
375    pub fn map(&self, logical_offset: u64) -> Option<(u64, u32)> {
376        let mut current_logical_start = 0;
377        for mapping in &self.mappings {
378            let current_logical_end = current_logical_start + mapping.length;
379            if logical_offset < current_logical_end {
380                let delta = logical_offset - current_logical_start;
381                let dev_offset = mapping.target_block_offset + delta;
382                let len = u32::try_from(current_logical_end - logical_offset).unwrap_or(u32::MAX);
383                return Some((dev_offset, len));
384            }
385            current_logical_start = current_logical_end;
386        }
387        None
388    }
389
390    pub fn is_empty(&self) -> bool {
391        self.mappings.is_empty()
392    }
393
394    pub fn total_blocks(&self) -> u64 {
395        self.mappings.iter().map(|m| m.length).sum()
396    }
397
398    /// Returns true if the block range `[offset, offset + length)` falls entirely within this
399    /// map's total blocks. If the map is empty, returns true (no range restriction).
400    pub fn are_blocks_within_source_range(&self, (offset, length): (u64, u32)) -> bool {
401        if self.is_empty() {
402            return true;
403        }
404        let total = self.total_blocks();
405        offset <= total && total - offset >= length as u64
406    }
407
408    pub fn mappings(&self) -> &[BlockOffsetMapping] {
409        &self.mappings
410    }
411}
412
413impl TryFrom<&fblock::BlockOffsetMapping> for BlockOffsetMapping {
414    type Error = zx::Status;
415
416    fn try_from(wire: &fblock::BlockOffsetMapping) -> Result<Self, Self::Error> {
417        if wire.length == 0 {
418            return Err(zx::Status::INVALID_ARGS);
419        }
420        wire.target_block_offset.checked_add(wire.length).ok_or(zx::Status::OUT_OF_RANGE)?;
421        Ok(BlockOffsetMapping {
422            target_block_offset: wire.target_block_offset,
423            length: wire.length,
424        })
425    }
426}
427
428impl TryFrom<fblock::BlockOffsetMapping> for BlockOffsetMapping {
429    type Error = zx::Status;
430
431    fn try_from(wire: fblock::BlockOffsetMapping) -> Result<Self, Self::Error> {
432        BlockOffsetMapping::try_from(&wire)
433    }
434}
435
436impl From<&BlockOffsetMapping> for fblock::BlockOffsetMapping {
437    fn from(m: &BlockOffsetMapping) -> Self {
438        fblock::BlockOffsetMapping { target_block_offset: m.target_block_offset, length: m.length }
439    }
440}
441
442impl From<BlockOffsetMapping> for fblock::BlockOffsetMapping {
443    fn from(m: BlockOffsetMapping) -> Self {
444        fblock::BlockOffsetMapping::from(&m)
445    }
446}
447
448impl TryFrom<&[fblock::BlockOffsetMapping]> for OffsetMap {
449    type Error = zx::Status;
450
451    fn try_from(wire: &[fblock::BlockOffsetMapping]) -> Result<Self, Self::Error> {
452        let raw_mappings: Vec<BlockOffsetMapping> =
453            wire.iter().map(BlockOffsetMapping::try_from).collect::<Result<_, _>>()?;
454        OffsetMap::new(raw_mappings)
455    }
456}
457
458impl TryFrom<Vec<fblock::BlockOffsetMapping>> for OffsetMap {
459    type Error = zx::Status;
460
461    fn try_from(wire: Vec<fblock::BlockOffsetMapping>) -> Result<Self, Self::Error> {
462        OffsetMap::try_from(wire.as_slice())
463    }
464}
465
466impl From<&OffsetMap> for Vec<fblock::BlockOffsetMapping> {
467    fn from(offset_map: &OffsetMap) -> Self {
468        offset_map.mappings.iter().map(fblock::BlockOffsetMapping::from).collect()
469    }
470}
471
472// Methods take Arc<Self> rather than &self because of
473// https://github.com/rust-lang/rust/issues/42940.
474pub trait SessionManager: 'static {
475    /// The Orchestrator is an object that holds the `SessionManager` and any other state that needs
476    /// to be shared between sessions.  It is responsible for keeping the `SessionManager` alive.
477    /// We use this type instead of directly holding an Arc<SessionManager> in BlockServer, to avoid
478    /// nested Arcs in concrete implementations which need to keep additional state.
479    type Orchestrator: Borrow<Self> + Send + Sync;
480
481    const SUPPORTS_DECOMPRESSION: bool;
482
483    type Session;
484
485    /// Returns true iff `a` and `b` identify the same session.  Used to scope
486    /// group-ID lookups in the shared `active_requests` slab to the originating
487    /// session.
488    fn session_eq(a: &Self::Session, b: &Self::Session) -> bool;
489
490    fn on_attach_vmo(
491        orchestrator: Arc<Self::Orchestrator>,
492        vmo: &Arc<zx::Vmo>,
493    ) -> impl Future<Output = Result<(), zx::Status>> + Send;
494
495    /// Creates a new session to handle `stream`.
496    ///
497    /// The returned future should run until the session completes, for example when the client end
498    /// closes.
499    ///
500    /// `offset_map` is an optional client-provided map to adjust the offset/length of FIFO
501    /// requests.  If the implementation supports mapping requests, it must forward this back to
502    /// [`SessionHelper::new`].
503    fn open_session(
504        orchestrator: Arc<Self::Orchestrator>,
505        stream: fblock::SessionRequestStream,
506        offset_map: OffsetMap,
507        block_size: u32,
508    ) -> impl Future<Output = Result<(), Error>> + Send;
509
510    /// Called to get block/partition information for Block::GetInfo, Partition::GetTypeGuid, etc.
511    fn get_info(&self) -> Cow<'_, DeviceInfo>;
512
513    /// Called to handle the GetVolumeInfo FIDL call.
514    fn get_volume_info(
515        &self,
516    ) -> impl Future<Output = Result<(fblock::VolumeManagerInfo, fblock::VolumeInfo), zx::Status>> + Send
517    {
518        async { Err(zx::Status::NOT_SUPPORTED) }
519    }
520
521    /// Called to handle the QuerySlices FIDL call.
522    fn query_slices(
523        &self,
524        _start_slices: &[u64],
525    ) -> impl Future<Output = Result<Vec<fblock::VsliceRange>, zx::Status>> + Send {
526        async { Err(zx::Status::NOT_SUPPORTED) }
527    }
528
529    /// Called to handle the Shrink FIDL call.
530    fn extend(
531        &self,
532        _start_slice: u64,
533        _slice_count: u64,
534    ) -> impl Future<Output = Result<(), zx::Status>> + Send {
535        async { Err(zx::Status::NOT_SUPPORTED) }
536    }
537
538    /// Called to handle the Shrink FIDL call.
539    fn shrink(
540        &self,
541        _start_slice: u64,
542        _slice_count: u64,
543    ) -> impl Future<Output = Result<(), zx::Status>> + Send {
544        async { Err(zx::Status::NOT_SUPPORTED) }
545    }
546
547    /// Opens a new mapper session.
548    fn open_mapper_session(
549        _orchestrator: Arc<Self::Orchestrator>,
550        _session: fidl::endpoints::ServerEnd<fblock::MapperSessionMarker>,
551        _mapping_vmo: zx::Vmo,
552        _block_size: u32,
553        _port: Option<zx::Port>,
554        _delivery_queue: Option<zx::Vmo>,
555    ) -> Result<impl Future<Output = Result<(), Error>> + Send, zx::Status> {
556        Err::<std::future::Ready<Result<(), Error>>, _>(zx::Status::NOT_SUPPORTED)
557    }
558
559    /// Returns the active requests.
560    fn active_requests(&self) -> &ActiveRequests<Self::Session>;
561}
562
563/// A helper trait for converting various types into an `Orchestrator`.
564///
565/// This exists to simplify [`BlockServer::new`].
566pub trait IntoOrchestrator {
567    type SM: SessionManager;
568
569    fn into_orchestrator(self) -> Arc<<Self::SM as SessionManager>::Orchestrator>;
570}
571
572impl<SM: SessionManager> BlockServer<SM> {
573    pub fn new(block_size: u32, orchestrator: impl IntoOrchestrator<SM = SM>) -> Self {
574        Self { block_size, orchestrator: orchestrator.into_orchestrator() }
575    }
576
577    pub fn session_manager(&self) -> &SM {
578        self.orchestrator.as_ref().borrow()
579    }
580
581    /// Called to process requests for fuchsia.storage.block.Block.
582    pub async fn handle_requests(
583        &self,
584        mut requests: fblock::BlockRequestStream,
585    ) -> Result<(), Error> {
586        let scope = fasync::Scope::new();
587        loop {
588            match requests.try_next().await {
589                Ok(Some(request)) => {
590                    if let Some(session) = self.handle_request(request, &scope).await? {
591                        scope.spawn(session.map(|_| ()));
592                    }
593                }
594                Ok(None) => break,
595                Err(error) => log::warn!(error:?; "Invalid request"),
596            }
597        }
598        scope.await;
599        Ok(())
600    }
601
602    /// Called to process requests for fuchsia.storage.block.Mapper.
603    pub async fn handle_mapper_requests(
604        &self,
605        requests: fblock::MapperRequestStream,
606    ) -> Result<(), Error> {
607        Self::handle_mapper_requests_impl(self.orchestrator.clone(), self.block_size, requests)
608            .await
609    }
610
611    async fn handle_mapper_requests_impl(
612        orchestrator: Arc<SM::Orchestrator>,
613        block_size: u32,
614        mut requests: fblock::MapperRequestStream,
615    ) -> Result<(), Error> {
616        let scope = fasync::Scope::new();
617        loop {
618            match requests.try_next().await {
619                Ok(Some(request)) => match request {
620                    fblock::MapperRequest::OpenSession {
621                        session,
622                        mapping_vmo,
623                        port,
624                        delivery_queue,
625                        responder,
626                    } => {
627                        match SM::open_mapper_session(
628                            orchestrator.clone(),
629                            session,
630                            mapping_vmo,
631                            block_size,
632                            port,
633                            delivery_queue,
634                        ) {
635                            Ok(fut) => {
636                                responder.send(Ok(()))?;
637                                scope.spawn(async move {
638                                    if let Err(error) = fut.await {
639                                        log::warn!(error:?; "Mapper session failed");
640                                    }
641                                });
642                            }
643                            Err(status) => {
644                                responder.send(Err(status.into_raw()))?;
645                            }
646                        }
647                    }
648                    fblock::MapperRequest::_UnknownMethod { .. } => {}
649                },
650                Ok(None) => break,
651                Err(error) => log::warn!(error:?; "Invalid mapper request"),
652            }
653        }
654        scope.await;
655        Ok(())
656    }
657
658    /// Processes a Block request.  If a new session task is created in response to the request,
659    /// it is returned.
660    async fn handle_request(
661        &self,
662        request: fblock::BlockRequest,
663        scope: &fasync::Scope,
664    ) -> Result<Option<impl Future<Output = Result<(), Error>> + Send + use<SM>>, Error> {
665        match request {
666            fblock::BlockRequest::GetInfo { responder } => {
667                let info = self.device_info();
668                let max_transfer_size = info.max_transfer_size(self.block_size);
669                let (block_count, mut flags) = match info.as_ref() {
670                    DeviceInfo::Block(BlockInfo { block_count, device_flags, .. }) => {
671                        (*block_count, *device_flags)
672                    }
673                    DeviceInfo::Partition(partition_info) => {
674                        (partition_info.block_count, partition_info.device_flags)
675                    }
676                    DeviceInfo::Volume(volume_info) => {
677                        let volume_info_fidl = self.session_manager().get_volume_info().await?;
678                        let block_count = volume_info_fidl.0.slice_size
679                            * volume_info_fidl.1.partition_slice_count
680                            / self.block_size as u64;
681                        (block_count, volume_info.device_flags)
682                    }
683                };
684                if SM::SUPPORTS_DECOMPRESSION {
685                    flags |= fblock::DeviceFlag::ZSTD_DECOMPRESSION_SUPPORT;
686                }
687                responder.send(Ok(&fblock::BlockInfo {
688                    block_count,
689                    block_size: self.block_size,
690                    max_transfer_size,
691                    flags,
692                }))?;
693            }
694            fblock::BlockRequest::OpenSession { session, control_handle: _ } => {
695                return Ok(Some(SM::open_session(
696                    self.orchestrator.clone(),
697                    session.into_stream(),
698                    OffsetMap::empty(),
699                    self.block_size,
700                )));
701            }
702            fblock::BlockRequest::OpenSessionWithOptions {
703                session,
704                mappings,
705                control_handle: _,
706            } => {
707                let info = self.device_info();
708                let offset_map: OffsetMap = match mappings.as_slice().try_into() {
709                    Ok(map) => map,
710                    Err(status) => {
711                        session.close_with_epitaph(status)?;
712                        return Ok(None);
713                    }
714                };
715                if let Some(max) = info.block_count() {
716                    for m in offset_map.mappings() {
717                        if m.target_block_offset.checked_add(m.length).unwrap_or(u64::MAX) > max {
718                            log::warn!("Invalid mapping for session: {m:?} (max blocks {max})");
719                            session.close_with_epitaph(zx::Status::OUT_OF_RANGE)?;
720                            return Ok(None);
721                        }
722                    }
723                }
724                return Ok(Some(SM::open_session(
725                    self.orchestrator.clone(),
726                    session.into_stream(),
727                    offset_map,
728                    self.block_size,
729                )));
730            }
731            fblock::BlockRequest::ConnectMapper { server_end, responder } => {
732                let orchestrator = self.orchestrator.clone();
733                let block_size = self.block_size;
734                scope.spawn(async move {
735                    if let Err(e) = Self::handle_mapper_requests_impl(
736                        orchestrator,
737                        block_size,
738                        server_end.into_stream(),
739                    )
740                    .await
741                    {
742                        log::warn!(e:?; "Error serving mapper requests");
743                    }
744                });
745                let _ = responder.send(Ok(()));
746            }
747            fblock::BlockRequest::GetTypeGuid { responder } => {
748                match self.device_info().type_guid() {
749                    Some(guid) => {
750                        responder.send(zx::sys::ZX_OK, Some(&fblock::Guid { value: guid }))?
751                    }
752                    None => responder.send(zx::sys::ZX_ERR_NOT_SUPPORTED, None)?,
753                }
754            }
755            fblock::BlockRequest::GetInstanceGuid { responder } => {
756                match self.device_info().instance_guid() {
757                    Some(guid) => {
758                        responder.send(zx::sys::ZX_OK, Some(&fblock::Guid { value: guid }))?
759                    }
760                    None => responder.send(zx::sys::ZX_ERR_NOT_SUPPORTED, None)?,
761                }
762            }
763            fblock::BlockRequest::GetName { responder } => {
764                let info = self.device_info();
765                match info.as_ref() {
766                    DeviceInfo::Partition(_) | DeviceInfo::Volume(_) => {
767                        responder.send(zx::sys::ZX_OK, Some(info.label()))?;
768                    }
769                    _ => responder.send(zx::sys::ZX_ERR_NOT_SUPPORTED, None)?,
770                }
771            }
772            fblock::BlockRequest::GetMetadata { responder } => {
773                let device_info = self.device_info();
774                match device_info.as_ref() {
775                    DeviceInfo::Partition(info) => {
776                        let mut type_guid =
777                            fblock::Guid { value: [0u8; fblock::GUID_LENGTH as usize] };
778                        type_guid.value.copy_from_slice(&info.type_guid);
779                        let mut instance_guid =
780                            fblock::Guid { value: [0u8; fblock::GUID_LENGTH as usize] };
781                        instance_guid.value.copy_from_slice(&info.instance_guid);
782                        let start_block_offset = info.start_block_offset;
783                        let flags = info.flags;
784                        responder.send(Ok(&fblock::PartitionInfo {
785                            name: Some(info.name.clone()),
786                            type_guid: Some(type_guid),
787                            instance_guid: Some(instance_guid),
788                            start_block_offset,
789                            num_blocks: device_info.block_count(),
790                            flags,
791                            ..Default::default()
792                        }))?;
793                    }
794                    DeviceInfo::Volume(info) => {
795                        let mut type_guid =
796                            fblock::Guid { value: [0u8; fblock::GUID_LENGTH as usize] };
797                        type_guid.value.copy_from_slice(&info.type_guid);
798                        let mut instance_guid =
799                            fblock::Guid { value: [0u8; fblock::GUID_LENGTH as usize] };
800                        instance_guid.value.copy_from_slice(&info.instance_guid);
801                        responder.send(Ok(&fblock::PartitionInfo {
802                            name: Some(info.name.clone()),
803                            type_guid: Some(type_guid),
804                            instance_guid: Some(instance_guid),
805                            start_block_offset: None,
806                            num_blocks: device_info.block_count(),
807                            flags: Some(info.flags),
808                            ..Default::default()
809                        }))?;
810                    }
811                    _ => responder.send(Err(zx::sys::ZX_ERR_NOT_SUPPORTED))?,
812                }
813            }
814            fblock::BlockRequest::QuerySlices { responder, start_slices } => {
815                match self.session_manager().query_slices(&start_slices).await {
816                    Ok(mut results) => {
817                        let results_len = results.len();
818                        assert!(results_len <= 16);
819                        results.resize(16, fblock::VsliceRange { allocated: false, count: 0 });
820                        responder.send(
821                            zx::sys::ZX_OK,
822                            &results.try_into().unwrap(),
823                            results_len as u64,
824                        )?;
825                    }
826                    Err(s) => {
827                        responder.send(
828                            s.into_raw(),
829                            &[fblock::VsliceRange { allocated: false, count: 0 }; 16],
830                            0,
831                        )?;
832                    }
833                }
834            }
835            fblock::BlockRequest::GetVolumeInfo { responder, .. } => {
836                match self.session_manager().get_volume_info().await {
837                    Ok((manager_info, volume_info)) => {
838                        responder.send(zx::sys::ZX_OK, Some(&manager_info), Some(&volume_info))?
839                    }
840                    Err(s) => responder.send(s.into_raw(), None, None)?,
841                }
842            }
843            fblock::BlockRequest::Extend { responder, start_slice, slice_count } => {
844                responder.send(zx::Status::result_into_raw(
845                    self.session_manager().extend(start_slice, slice_count).await,
846                ))?;
847            }
848            fblock::BlockRequest::Shrink { responder, start_slice, slice_count } => {
849                responder.send(zx::Status::result_into_raw(
850                    self.session_manager().shrink(start_slice, slice_count).await,
851                ))?;
852            }
853            fblock::BlockRequest::Destroy { responder, .. } => {
854                responder.send(zx::sys::ZX_ERR_NOT_SUPPORTED)?;
855            }
856        }
857        Ok(None)
858    }
859
860    fn device_info(&self) -> Cow<'_, DeviceInfo> {
861        self.session_manager().get_info()
862    }
863}
864
865pub(crate) struct RegisteredVmo {
866    pub vmo: Arc<zx::Vmo>,
867    pub size: u64,
868    pub mapping: Option<Arc<VmoMapping>>,
869}
870
871impl RegisteredVmo {
872    /// Validates that a request range (`vmo_offset` to `vmo_offset + length`, in bytes) falls
873    /// within the bounds of this VMO.
874    fn validate_request(&self, vmo_offset: u64, length: u64) -> Result<(), zx::Status> {
875        if vmo_offset > self.size || self.size - vmo_offset < length {
876            Err(zx::Status::OUT_OF_RANGE)
877        } else {
878            Ok(())
879        }
880    }
881
882    /// Returns the cached VMO mapping if available, or creates and caches a new mapping.
883    fn get_or_create_mapping(&mut self) -> Result<Arc<VmoMapping>, zx::Status> {
884        match &self.mapping {
885            Some(mapping) => Ok(mapping.clone()),
886            None => {
887                let mapping = VmoMapping::new(&self.vmo, self.size as usize)?;
888                self.mapping = Some(mapping.clone());
889                Ok(mapping)
890            }
891        }
892    }
893}
894
895struct SessionHelper<SM: SessionManager> {
896    orchestrator: Arc<SM::Orchestrator>,
897    offset_map: OffsetMap,
898    max_transfer_blocks: Option<NonZero<u32>>,
899    block_size: u32,
900    peer_fifo: zx::Fifo<BlockFifoResponse, BlockFifoRequest>,
901    vmos: Mutex<BTreeMap<u16, RegisteredVmo>>,
902}
903
904struct VmoMapping {
905    base: usize,
906    size: usize,
907}
908
909impl VmoMapping {
910    fn new(vmo: &zx::Vmo, size: usize) -> Result<Arc<Self>, zx::Status> {
911        Ok(Arc::new(Self {
912            base: fuchsia_runtime::vmar_root_self()
913                .map(0, vmo, 0, size, zx::VmarFlags::PERM_WRITE | zx::VmarFlags::PERM_READ)
914                .inspect_err(|error| {
915                    log::warn!(error:?, size; "VmoMapping: unable to map VMO");
916                })?,
917            size,
918        }))
919    }
920}
921
922impl Drop for VmoMapping {
923    fn drop(&mut self) {
924        // SAFETY: We mapped this in `VmoMapping::new`.
925        unsafe {
926            let _ = fuchsia_runtime::vmar_root_self().unmap(self.base, self.size);
927        }
928    }
929}
930
931enum HandleRequestResult {
932    /// The request was handled successfully.
933    Ok,
934    /// The request closed the stream.  The caller must shut down the session, and must call the
935    /// provided callback after the session is completely shut down.  The caller should assume that
936    /// no further requests need to be handled once this is received.
937    Closed(Box<dyn FnOnce() + Send + 'static>),
938}
939
940impl<SM: SessionManager> SessionHelper<SM> {
941    fn new(
942        orchestrator: Arc<SM::Orchestrator>,
943        offset_map: OffsetMap,
944        max_transfer_blocks: Option<NonZero<u32>>,
945        block_size: u32,
946    ) -> Result<(Self, zx::Fifo<BlockFifoRequest, BlockFifoResponse>), zx::Status> {
947        let (peer_fifo, fifo) = zx::Fifo::create(16)?;
948        Ok((
949            Self {
950                orchestrator,
951                offset_map,
952                max_transfer_blocks,
953                block_size,
954                peer_fifo,
955                vmos: Mutex::default(),
956            },
957            fifo,
958        ))
959    }
960
961    fn session_manager(&self) -> &SM {
962        self.orchestrator.as_ref().borrow()
963    }
964
965    async fn handle_request(
966        &self,
967        request: fblock::SessionRequest,
968    ) -> Result<HandleRequestResult, Error> {
969        match request {
970            fblock::SessionRequest::GetFifo { responder } => {
971                let rights = zx::Rights::TRANSFER
972                    | zx::Rights::READ
973                    | zx::Rights::WRITE
974                    | zx::Rights::SIGNAL
975                    | zx::Rights::WAIT;
976                match self.peer_fifo.duplicate_handle(rights) {
977                    Ok(fifo) => responder.send(Ok(fifo.downcast()))?,
978                    Err(s) => responder.send(Err(s.into_raw()))?,
979                }
980                Ok(HandleRequestResult::Ok)
981            }
982            fblock::SessionRequest::AttachVmo { vmo, responder } => {
983                let info = vmo.info().map_err(Error::from)?;
984                if info.flags.contains(zx::VmoInfoFlags::RESIZABLE) {
985                    responder.send(Err(zx::Status::INVALID_ARGS.into_raw()))?;
986                    return Ok(HandleRequestResult::Ok);
987                }
988                let size = info.size_bytes;
989                let vmo = Arc::new(vmo);
990                let vmo_id = {
991                    let mut vmos = self.vmos.lock();
992                    if vmos.len() == u16::MAX as usize {
993                        responder.send(Err(zx::Status::NO_RESOURCES.into_raw()))?;
994                        return Ok(HandleRequestResult::Ok);
995                    } else {
996                        let vmo_id = match vmos.last_entry() {
997                            None => 1,
998                            Some(o) => {
999                                o.key().checked_add(1).unwrap_or_else(|| {
1000                                    let mut vmo_id = 1;
1001                                    // Find the first gap...
1002                                    for (&id, _) in &*vmos {
1003                                        if id > vmo_id {
1004                                            break;
1005                                        }
1006                                        vmo_id = id + 1;
1007                                    }
1008                                    vmo_id
1009                                })
1010                            }
1011                        };
1012                        vmos.insert(
1013                            vmo_id,
1014                            RegisteredVmo { vmo: vmo.clone(), size, mapping: None },
1015                        );
1016                        vmo_id
1017                    }
1018                };
1019                SM::on_attach_vmo(self.orchestrator.clone(), &vmo).await?;
1020                responder.send(Ok(&fblock::VmoId { id: vmo_id }))?;
1021                Ok(HandleRequestResult::Ok)
1022            }
1023            fblock::SessionRequest::Close { responder } => {
1024                Ok(HandleRequestResult::Closed(Box::new(move || {
1025                    if let Err(error) = responder.send(Ok(())) {
1026                        log::warn!(error:?; "Error sending close response");
1027                    }
1028                })))
1029            }
1030        }
1031    }
1032
1033    /// Decodes `request`.
1034    fn decode_fifo_request(
1035        &self,
1036        session: SM::Session,
1037        request: &BlockFifoRequest,
1038    ) -> Result<DecodedRequest, Option<BlockFifoResponse>> {
1039        let flags = BlockIoFlag::from_bits_truncate(request.command.flags);
1040
1041        let request_bytes = request.length as u64 * self.block_size as u64;
1042
1043        let mut operation = BlockOpcode::from_primitive(request.command.opcode)
1044            .ok_or(zx::Status::INVALID_ARGS)
1045            .and_then(|code| {
1046                if flags.contains(BlockIoFlag::DECOMPRESS_WITH_ZSTD) {
1047                    if code != BlockOpcode::Read {
1048                        return Err(zx::Status::INVALID_ARGS);
1049                    }
1050                    if !SM::SUPPORTS_DECOMPRESSION {
1051                        return Err(zx::Status::NOT_SUPPORTED);
1052                    }
1053                }
1054                if matches!(code, BlockOpcode::Read | BlockOpcode::Write | BlockOpcode::Trim) {
1055                    if request.length == 0 {
1056                        return Err(zx::Status::INVALID_ARGS);
1057                    }
1058                    // Make sure the end offset won't wrap.
1059                    if request.dev_offset.checked_add(request.length as u64).is_none() {
1060                        return Err(zx::Status::OUT_OF_RANGE);
1061                    }
1062                }
1063                if matches!(code, BlockOpcode::Read | BlockOpcode::Write) {
1064                    let vmo_byte_offset = request
1065                        .vmo_offset
1066                        .checked_mul(self.block_size as u64)
1067                        .ok_or(zx::Status::OUT_OF_RANGE)?;
1068                    if request_bytes.checked_add(vmo_byte_offset).is_none() {
1069                        return Err(zx::Status::OUT_OF_RANGE);
1070                    }
1071                }
1072                Ok(match code {
1073                    BlockOpcode::Read => Operation::Read {
1074                        device_block_offset: request.dev_offset,
1075                        block_count: request.length,
1076                        _unused: 0,
1077                        vmo_offset: request
1078                            .vmo_offset
1079                            .checked_mul(self.block_size as u64)
1080                            .ok_or(zx::Status::OUT_OF_RANGE)?,
1081                        options: ReadOptions {
1082                            inline_crypto: InlineCryptoOptions {
1083                                is_enabled: flags.contains(BlockIoFlag::INLINE_ENCRYPTION_ENABLED),
1084                                dun: request.dun,
1085                                slot: request.slot,
1086                            },
1087                        },
1088                    },
1089                    BlockOpcode::Write => {
1090                        let mut options = WriteOptions {
1091                            inline_crypto: InlineCryptoOptions {
1092                                is_enabled: flags.contains(BlockIoFlag::INLINE_ENCRYPTION_ENABLED),
1093                                dun: request.dun,
1094                                slot: request.slot,
1095                            },
1096                            ..WriteOptions::default()
1097                        };
1098                        if flags.contains(BlockIoFlag::FORCE_ACCESS) {
1099                            options.flags |= WriteFlags::FORCE_ACCESS;
1100                        }
1101                        if flags.contains(BlockIoFlag::PRE_BARRIER) {
1102                            options.flags |= WriteFlags::PRE_BARRIER;
1103                        }
1104                        Operation::Write {
1105                            device_block_offset: request.dev_offset,
1106                            block_count: request.length,
1107                            _unused: 0,
1108                            options,
1109                            vmo_offset: request
1110                                .vmo_offset
1111                                .checked_mul(self.block_size as u64)
1112                                .ok_or(zx::Status::OUT_OF_RANGE)?,
1113                        }
1114                    }
1115                    BlockOpcode::Flush => Operation::Flush,
1116                    BlockOpcode::Trim => Operation::Trim {
1117                        device_block_offset: request.dev_offset,
1118                        block_count: request.length,
1119                    },
1120                    BlockOpcode::CloseVmo => Operation::CloseVmo,
1121                })
1122            });
1123
1124        let group_or_request = if flags.contains(BlockIoFlag::GROUP_ITEM) {
1125            GroupOrRequest::Group(request.group)
1126        } else {
1127            GroupOrRequest::Request(request.reqid)
1128        };
1129
1130        let mut active_requests = self.session_manager().active_requests().0.lock();
1131        let mut request_id = None;
1132
1133        // Multiple Block I/O request may be sent as a group.
1134        // Notes:
1135        // - the group is identified by the group id in the request
1136        // - if using groups, a response will not be sent unless `BlockIoFlag::GROUP_LAST`
1137        //   flag is set.
1138        // - when processing a request of a group fails, subsequent requests of that
1139        //   group will not be processed.
1140        // - decompression is a special case, see block-fifo.h for semantics.
1141        //
1142        // Refer to sdk/fidl/fuchsia.hardware.block.driver/block.fidl for details.
1143        if group_or_request.is_group() {
1144            // Search for an existing entry that matches this group.  NOTE: This is a potentially
1145            // expensive way to find a group (it's iterating over all slots in the active-requests
1146            // slab).  This can be optimised easily should we need to.
1147            for (key, group) in &mut active_requests.requests {
1148                if group.group_or_request == group_or_request
1149                    && SM::session_eq(&group.session, &session)
1150                {
1151                    if group.req_id.is_some() {
1152                        // We have already received a request tagged as last.
1153                        if group.status.is_ok() {
1154                            group.status = Err(zx::Status::INVALID_ARGS);
1155                        }
1156                        // Ignore this request.
1157                        return Err(None);
1158                    }
1159                    // See if this is a continuation of a decompressed read.
1160                    if group.status.is_ok()
1161                        && let Some(info) = &mut group.decompression_info
1162                    {
1163                        if let Ok(Operation::Read {
1164                            device_block_offset,
1165                            mut block_count,
1166                            options,
1167                            vmo_offset: 0,
1168                            ..
1169                        }) = operation
1170                        {
1171                            let remaining_bytes = info
1172                                .compressed_range
1173                                .end
1174                                .next_multiple_of(self.block_size as usize)
1175                                as u64
1176                                - info.bytes_so_far;
1177                            if !flags.contains(BlockIoFlag::DECOMPRESS_WITH_ZSTD)
1178                                || request.total_compressed_bytes != 0
1179                                || request.uncompressed_bytes != 0
1180                                || request.compressed_prefix_bytes != 0
1181                                || (flags.contains(BlockIoFlag::GROUP_LAST)
1182                                    && info.bytes_so_far + request_bytes
1183                                        < info.compressed_range.end as u64)
1184                                || (!flags.contains(BlockIoFlag::GROUP_LAST)
1185                                    && request_bytes >= remaining_bytes)
1186                            {
1187                                group.status = Err(zx::Status::INVALID_ARGS);
1188                            } else {
1189                                // We are tolerant of `block_count` being more than we actually
1190                                // need.  This can happen if the client is working with a larger
1191                                // block size than the device block size.  For example, if Blobfs
1192                                // has a 8192 byte block size, but the device might has a 512 byte
1193                                // block size, it can ask for a multiple of 16 blocks, when fewer
1194                                // than that might actually be required to hold the compressed data.
1195                                // It is easier for us to tolerate this here than to get Blobfs to
1196                                // change to pass only the blocks that are required.
1197                                if request_bytes > remaining_bytes {
1198                                    block_count = (remaining_bytes / self.block_size as u64) as u32;
1199                                }
1200
1201                                operation = Ok(Operation::ContinueDecompressedRead {
1202                                    offset: info.bytes_so_far,
1203                                    device_block_offset,
1204                                    block_count,
1205                                    options,
1206                                });
1207
1208                                info.bytes_so_far += block_count as u64 * self.block_size as u64;
1209                            }
1210                        } else {
1211                            group.status = Err(zx::Status::INVALID_ARGS);
1212                        }
1213                    }
1214                    if flags.contains(BlockIoFlag::GROUP_LAST) {
1215                        group.req_id = Some(request.reqid);
1216                        // If the group has had an error, there is no point trying to issue this
1217                        // request.
1218                        if let Err(s) = group.status {
1219                            operation = Err(s);
1220                        }
1221                    } else if group.status.is_err() {
1222                        // The group has already encountered an error, so there is no point trying
1223                        // to issue this request.
1224                        return Err(None);
1225                    }
1226                    request_id = Some(RequestId(key));
1227                    group.count += 1;
1228                    break;
1229                }
1230            }
1231        }
1232
1233        let is_single_request =
1234            !flags.contains(BlockIoFlag::GROUP_ITEM) || flags.contains(BlockIoFlag::GROUP_LAST);
1235
1236        let mut decompression_info = None;
1237        let vmo = match operation {
1238            Ok(Operation::Read {
1239                device_block_offset,
1240                mut block_count,
1241                options,
1242                vmo_offset,
1243                ..
1244            }) => match self.vmos.lock().get_mut(&request.vmoid) {
1245                Some(registered_vmo) => {
1246                    if flags.contains(BlockIoFlag::DECOMPRESS_WITH_ZSTD) {
1247                        let compressed_range = request.compressed_prefix_bytes as usize
1248                            ..request.compressed_prefix_bytes as usize
1249                                + request.total_compressed_bytes as usize;
1250                        let required_buffer_size =
1251                            compressed_range.end.next_multiple_of(self.block_size as usize);
1252
1253                        // Validate the initial decompression request.
1254                        if compressed_range.start >= compressed_range.end
1255                            || vmo_offset.checked_add(request.uncompressed_bytes as u64).is_none()
1256                            || (is_single_request && request_bytes < compressed_range.end as u64)
1257                            || (!is_single_request && request_bytes >= required_buffer_size as u64)
1258                        {
1259                            Err(zx::Status::INVALID_ARGS)
1260                        } else {
1261                            // We are tolerant of `block_count` being more than we actually need.
1262                            // This can happen if the client is working in a larger block size than
1263                            // the device block size.  For example, Blobfs has a 8192 byte block
1264                            // size, but the device might have a 512 byte block size.  It is easier
1265                            // for us to tolerate this here than to get Blobfs to change to pass
1266                            // only the blocks that are required.
1267                            let bytes_so_far = if request_bytes > required_buffer_size as u64 {
1268                                block_count =
1269                                    (required_buffer_size / self.block_size as usize) as u32;
1270                                required_buffer_size as u64
1271                            } else {
1272                                request_bytes
1273                            };
1274
1275                            // To decompress, we need to have the target VMO mapped (cached).
1276                            registered_vmo
1277                                .get_or_create_mapping()
1278                                .and_then(|mapping| {
1279                                    // Make sure `vmo_offset` and `uncompressed_bytes` are
1280                                    // within range.
1281                                    if vmo_offset
1282                                        .checked_add(request.uncompressed_bytes as u64)
1283                                        .is_some_and(|end| end <= mapping.size as u64)
1284                                    {
1285                                        Ok(mapping)
1286                                    } else {
1287                                        Err(zx::Status::OUT_OF_RANGE)
1288                                    }
1289                                })
1290                                .map(|mapping| {
1291                                    // Convert the operation into a `StartDecompressedRead`
1292                                    // operation. For non-fragmented requests, this will be the only
1293                                    // operation, but if it's a fragmented read,
1294                                    // `ContinueDecompressedRead` operations will follow.
1295                                    operation = Ok(Operation::StartDecompressedRead {
1296                                        required_buffer_size,
1297                                        device_block_offset,
1298                                        block_count,
1299                                        options,
1300                                    });
1301                                    // Record sufficient information so that we can decompress when
1302                                    // all the requests complete.
1303                                    decompression_info = Some(DecompressionInfo {
1304                                        compressed_range,
1305                                        bytes_so_far,
1306                                        mapping,
1307                                        uncompressed_range: vmo_offset
1308                                            ..vmo_offset + request.uncompressed_bytes as u64,
1309                                        buffer: None,
1310                                    });
1311                                    None
1312                                })
1313                        }
1314                    } else {
1315                        registered_vmo
1316                            .validate_request(vmo_offset, request_bytes)
1317                            .map(|()| Some(registered_vmo.vmo.clone()))
1318                    }
1319                }
1320                None => Err(zx::Status::IO),
1321            },
1322            Ok(Operation::Write { vmo_offset, .. }) => {
1323                self.vmos.lock().get(&request.vmoid).map_or(Err(zx::Status::IO), |registered_vmo| {
1324                    registered_vmo
1325                        .validate_request(vmo_offset, request_bytes)
1326                        .map(|()| Some(registered_vmo.vmo.clone()))
1327                })
1328            }
1329            Ok(Operation::CloseVmo) => {
1330                self.vmos.lock().remove(&request.vmoid).map_or(
1331                    Err(zx::Status::IO),
1332                    |registered_vmo| {
1333                        let vmo_clone = registered_vmo.vmo.clone();
1334                        // Make sure the VMO is dropped after all current Epoch guards have been
1335                        // dropped.
1336                        Epoch::global().defer(move || drop(vmo_clone));
1337                        Ok(Some(registered_vmo.vmo))
1338                    },
1339                )
1340            }
1341            _ => Ok(None),
1342        }
1343        .unwrap_or_else(|e| {
1344            operation = Err(e);
1345            None
1346        });
1347
1348        let trace_flow_id = NonZero::new(request.trace_flow_id);
1349        let request_id = request_id.unwrap_or_else(|| {
1350            RequestId(active_requests.requests.insert(ActiveRequest {
1351                session,
1352                group_or_request,
1353                trace_flow_id,
1354                _epoch_guard: Epoch::global().guard(),
1355                status: Ok(()),
1356                count: 1,
1357                req_id: is_single_request.then_some(request.reqid),
1358                decompression_info,
1359            }))
1360        });
1361
1362        Ok(DecodedRequest {
1363            request_id,
1364            trace_flow_id,
1365            operation: operation.map_err(|status| {
1366                active_requests.complete_and_take_response(request_id, Err(status)).map(|(_, r)| r)
1367            })?,
1368            vmo,
1369        })
1370    }
1371
1372    fn take_vmos(&self) -> BTreeMap<u16, RegisteredVmo> {
1373        std::mem::take(&mut *self.vmos.lock())
1374    }
1375
1376    /// Maps the request and returns the mapped request with an optional remainder.
1377    fn map_request(
1378        &self,
1379        mut request: DecodedRequest,
1380        active_request: &mut ActiveRequest<SM::Session>,
1381    ) -> Result<(DecodedRequest, Option<DecodedRequest>), zx::Status> {
1382        if active_request.status.is_err() {
1383            return Err(zx::Status::BAD_STATE);
1384        }
1385        if let Some(blocks) = request.operation.blocks() {
1386            if !self.offset_map.are_blocks_within_source_range(blocks) {
1387                return Err(zx::Status::OUT_OF_RANGE);
1388            }
1389        }
1390        let remainder =
1391            request.operation.map(&self.offset_map, self.max_transfer_blocks, self.block_size)?;
1392        if remainder.is_some() {
1393            active_request.count += 1;
1394        }
1395        static CACHE: AtomicU64 = AtomicU64::new(0);
1396        if let Some(context) =
1397            fuchsia_trace::TraceCategoryContext::acquire_cached("storage", &CACHE)
1398        {
1399            use fuchsia_trace::ArgValue;
1400            let trace_args = [
1401                ArgValue::of("request_id", request.request_id.0),
1402                ArgValue::of("opcode", request.operation.trace_label()),
1403            ];
1404            let _scope =
1405                fuchsia_trace::duration("storage", "block_server::start_transaction", &trace_args);
1406            if let Some(trace_flow_id) = active_request.trace_flow_id {
1407                fuchsia_trace::flow_step(
1408                    &context,
1409                    "block_server::start_transaction",
1410                    trace_flow_id.get().into(),
1411                    &[],
1412                );
1413            }
1414        }
1415        let remainder = remainder.map(|operation| DecodedRequest { operation, ..request.clone() });
1416        Ok((request, remainder))
1417    }
1418
1419    /// Drops all requests for which `pred` is true.
1420    ///
1421    /// NOTE: This should only be called once we are certain that the requests will not be
1422    /// completed asynchronously  Otherwise, requests might be completed twice.
1423    fn drop_active_requests(&self, pred: impl Fn(&SM::Session) -> bool) {
1424        self.session_manager().active_requests().0.lock().requests.retain(|_, r| !pred(&r.session));
1425    }
1426
1427    /// Closes all grouped requests for which `pred` is true and which are held open pending the
1428    /// completion of their group.
1429    ///
1430    /// Normally, a request is dropped from ActiveRequests when it is completed.  However, if a
1431    /// request is part of a group, it will not be dropped until a request with GROUP_LAST arrives.
1432    /// If we're shutting down a session, the client may not ever send the GROUP_LAST, so we need to
1433    /// be sure to close these grouped requests.
1434    ///
1435    /// This is called during session shutdown in situations where [`Self::drop_active_requests`]
1436    /// cannot be used (e.g. for the callback interface, which hands off the responsibility of
1437    /// completing requests to its concrete implementation and cannot control when requests are
1438    /// completed relative to session shutdown).
1439    fn close_active_groups(&self, pred: impl Fn(&SM::Session) -> bool) {
1440        self.session_manager().active_requests().0.lock().requests.retain(|_, request| {
1441            if !pred(&request.session) || request.req_id.is_some() {
1442                return true;
1443            }
1444            // Mark the group as completed, and immediately drop any which have no outstanding
1445            // requests (since they will otherwise never be dropped).
1446            request.req_id = Some(u32::MAX);
1447            request.count > 0
1448        });
1449    }
1450}
1451
1452#[repr(transparent)]
1453#[derive(Clone, Copy, Debug, Eq, PartialEq, Hash, Ord, PartialOrd)]
1454pub struct RequestId(usize);
1455
1456#[derive(Clone, Debug)]
1457struct DecodedRequest {
1458    request_id: RequestId,
1459    trace_flow_id: TraceFlowId,
1460    operation: Operation,
1461    vmo: Option<Arc<zx::Vmo>>,
1462}
1463
1464/// cbindgen:no-export
1465pub type WriteFlags = block_protocol::WriteFlags;
1466pub type WriteOptions = block_protocol::WriteOptions;
1467pub type ReadOptions = block_protocol::ReadOptions;
1468pub type InlineCryptoOptions = block_protocol::InlineCryptoOptions;
1469
1470#[repr(C)]
1471#[derive(Clone, Debug, PartialEq, Eq)]
1472pub enum Operation {
1473    // NOTE: On the C++ side, this ends up as a union and, for efficiency reasons, there is code
1474    // that assumes that some fields for reads and writes (and possibly trim) line-up (e.g. common
1475    // code can read `device_block_offset` from the read variant and then assume it's valid for the
1476    // write variant).
1477    Read {
1478        device_block_offset: u64,
1479        block_count: u32,
1480        _unused: u32,
1481        vmo_offset: u64,
1482        options: ReadOptions,
1483    },
1484    Write {
1485        device_block_offset: u64,
1486        block_count: u32,
1487        _unused: u32,
1488        vmo_offset: u64,
1489        options: WriteOptions,
1490    },
1491    Flush,
1492    Trim {
1493        device_block_offset: u64,
1494        block_count: u32,
1495    },
1496    /// This will never be seen by the C interface.
1497    CloseVmo,
1498    /// This will never be seen by the C interface.
1499    StartDecompressedRead {
1500        required_buffer_size: usize,
1501        device_block_offset: u64,
1502        block_count: u32,
1503        options: ReadOptions,
1504    },
1505    /// This will never be seen by the C interface.
1506    ContinueDecompressedRead {
1507        offset: u64,
1508        device_block_offset: u64,
1509        block_count: u32,
1510        options: ReadOptions,
1511    },
1512}
1513
1514impl Operation {
1515    fn trace_label(&self) -> &'static str {
1516        match self {
1517            Operation::Read { .. } => "read",
1518            Operation::Write { .. } => "write",
1519            Operation::Flush { .. } => "flush",
1520            Operation::Trim { .. } => "trim",
1521            Operation::CloseVmo { .. } => "close_vmo",
1522            Operation::StartDecompressedRead { .. } => "start_decompressed_read",
1523            Operation::ContinueDecompressedRead { .. } => "continue_decompressed_read",
1524        }
1525    }
1526
1527    /// Returns (offset, length).
1528    pub fn blocks(&self) -> Option<(u64, u32)> {
1529        match self {
1530            Operation::Read { device_block_offset, block_count, .. }
1531            | Operation::Write { device_block_offset, block_count, .. }
1532            | Operation::Trim { device_block_offset, block_count, .. } => {
1533                Some((*device_block_offset, *block_count))
1534            }
1535            _ => None,
1536        }
1537    }
1538
1539    /// Returns mutable references to (offset, length).
1540    fn blocks_mut(&mut self) -> Option<(&mut u64, &mut u32)> {
1541        match self {
1542            Operation::Read { device_block_offset, block_count, .. }
1543            | Operation::Write { device_block_offset, block_count, .. }
1544            | Operation::Trim { device_block_offset, block_count, .. } => {
1545                Some((device_block_offset, block_count))
1546            }
1547            _ => None,
1548        }
1549    }
1550
1551    /// Maps the operation using `offset_map` and returns the remainder if the request was split
1552    /// due to `max_transfer_blocks` or crossing mapping boundaries.
1553    fn map(
1554        &mut self,
1555        offset_map: &OffsetMap,
1556        max_transfer_blocks: Option<NonZero<u32>>,
1557        block_size: u32,
1558    ) -> Result<Option<Self>, zx::Status> {
1559        let mut max = match self {
1560            Operation::Read { .. } | Operation::Write { .. } => max_transfer_blocks.map(u32::from),
1561            _ => None,
1562        };
1563        let (offset, length) = match self.blocks_mut() {
1564            Some(b) => b,
1565            None => return Ok(None),
1566        };
1567        let orig_offset = *offset;
1568        if !offset_map.is_empty() {
1569            let (dev_offset, len) = offset_map.map(*offset).ok_or(zx::Status::OUT_OF_RANGE)?;
1570            *offset = dev_offset;
1571            max = match max {
1572                None => Some(len),
1573                Some(m) => Some(std::cmp::min(m, len)),
1574            };
1575        }
1576        if let Some(max) = max {
1577            if *length as u64 > max as u64 {
1578                let rem = *length - max;
1579                *length = max;
1580                return Ok(Some(match self {
1581                    Operation::Read {
1582                        device_block_offset: _,
1583                        block_count: _,
1584                        vmo_offset,
1585                        _unused,
1586                        options,
1587                    } => {
1588                        let mut options = *options;
1589                        options.inline_crypto.dun =
1590                            options.inline_crypto.dun.wrapping_add(max as u64);
1591                        Operation::Read {
1592                            device_block_offset: orig_offset + max as u64,
1593                            block_count: rem,
1594                            vmo_offset: *vmo_offset + max as u64 * block_size as u64,
1595                            _unused: *_unused,
1596                            options: options,
1597                        }
1598                    }
1599                    Operation::Write {
1600                        device_block_offset: _,
1601                        block_count: _,
1602                        _unused,
1603                        vmo_offset,
1604                        options,
1605                    } => {
1606                        let mut options = *options;
1607                        options.inline_crypto.dun =
1608                            options.inline_crypto.dun.wrapping_add(max as u64);
1609                        Operation::Write {
1610                            device_block_offset: orig_offset + max as u64,
1611                            block_count: rem,
1612                            _unused: *_unused,
1613                            vmo_offset: *vmo_offset + max as u64 * block_size as u64,
1614                            options: options,
1615                        }
1616                    }
1617                    Operation::Trim { device_block_offset: _, block_count: _ } => Operation::Trim {
1618                        device_block_offset: orig_offset + max as u64,
1619                        block_count: rem,
1620                    },
1621                    _ => unreachable!(),
1622                }));
1623            }
1624        }
1625        Ok(None)
1626    }
1627
1628    /// Returns true if the specified write flags are set.
1629    pub fn has_write_flag(&self, value: WriteFlags) -> bool {
1630        if let Operation::Write { options, .. } = self {
1631            options.flags.contains(value)
1632        } else {
1633            false
1634        }
1635    }
1636
1637    /// Removes `value` from the request's write flags and returns true if the flag was set.
1638    pub fn take_write_flag(&mut self, value: WriteFlags) -> bool {
1639        if let Operation::Write { options, .. } = self {
1640            let result = options.flags.contains(value);
1641            options.flags.remove(value);
1642            result
1643        } else {
1644            false
1645        }
1646    }
1647}
1648
1649#[derive(Clone, Copy, Debug, Eq, PartialEq, Ord, PartialOrd)]
1650pub enum GroupOrRequest {
1651    Group(u16),
1652    Request(u32),
1653}
1654
1655impl GroupOrRequest {
1656    fn is_group(&self) -> bool {
1657        matches!(self, Self::Group(_))
1658    }
1659
1660    fn group_id(&self) -> Option<u16> {
1661        match self {
1662            Self::Group(id) => Some(*id),
1663            Self::Request(_) => None,
1664        }
1665    }
1666}
1667
1668#[cfg(test)]
1669mod tests {
1670    use super::{
1671        BlockOffsetMapping, BlockServer, DeviceInfo, FIFO_MAX_REQUESTS, OffsetMap, Operation,
1672        PartitionInfo, TraceFlowId,
1673    };
1674    use assert_matches::assert_matches;
1675    use block_protocol::{
1676        BlockFifoCommand, BlockFifoRequest, BlockFifoResponse, InlineCryptoOptions, ReadOptions,
1677        WriteFlags, WriteOptions,
1678    };
1679    use fidl_fuchsia_storage_block as fblock;
1680    use fidl_fuchsia_storage_block::{BlockIoFlag, BlockOpcode};
1681    use fuchsia_async as fasync;
1682    use fuchsia_sync::Mutex;
1683    use futures::FutureExt as _;
1684    use futures::channel::oneshot;
1685    use futures::future::BoxFuture;
1686    use std::borrow::Cow;
1687    use std::future::poll_fn;
1688    use std::num::NonZero;
1689    use std::pin::pin;
1690    use std::sync::Arc;
1691    use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
1692    use std::task::{Context, Poll};
1693
1694    #[derive(Default)]
1695    struct MockInterface {
1696        info: Option<DeviceInfo>,
1697        attached_vmos: AtomicU64,
1698        read_hook: Option<
1699            Box<
1700                dyn Fn(u64, u32, &Arc<zx::Vmo>, u64) -> BoxFuture<'static, Result<(), zx::Status>>
1701                    + Send
1702                    + Sync,
1703            >,
1704        >,
1705        write_hook:
1706            Option<Box<dyn Fn(u64) -> BoxFuture<'static, Result<(), zx::Status>> + Send + Sync>>,
1707        barrier_hook: Option<Box<dyn Fn() -> Result<(), zx::Status> + Send + Sync>>,
1708    }
1709
1710    impl super::async_interface::Interface for MockInterface {
1711        async fn on_attach_vmo(&self, _vmo: &zx::Vmo) -> Result<(), zx::Status> {
1712            self.attached_vmos.fetch_add(1, Ordering::Relaxed);
1713            Ok(())
1714        }
1715
1716        fn on_detach_vmo(&self, _vmo: &zx::Vmo) {
1717            self.attached_vmos.fetch_sub(1, Ordering::Relaxed);
1718        }
1719
1720        fn get_info(&self) -> Cow<'_, DeviceInfo> {
1721            match &self.info {
1722                Some(info) => Cow::Borrowed(info),
1723                None => Cow::Owned(test_device_info()),
1724            }
1725        }
1726
1727        async fn read(
1728            &self,
1729            device_block_offset: u64,
1730            block_count: u32,
1731            vmo: &Arc<zx::Vmo>,
1732            vmo_offset: u64,
1733            _opts: ReadOptions,
1734            _trace_flow_id: TraceFlowId,
1735        ) -> Result<(), zx::Status> {
1736            if let Some(total) = self.get_info().block_count() {
1737                if device_block_offset >= total || total - device_block_offset < block_count as u64
1738                {
1739                    return Err(zx::Status::OUT_OF_RANGE);
1740                }
1741            }
1742            if let Some(read_hook) = &self.read_hook {
1743                read_hook(device_block_offset, block_count, vmo, vmo_offset).await
1744            } else {
1745                unimplemented!();
1746            }
1747        }
1748
1749        async fn write(
1750            &self,
1751            device_block_offset: u64,
1752            _block_count: u32,
1753            _vmo: &Arc<zx::Vmo>,
1754            _vmo_offset: u64,
1755            opts: WriteOptions,
1756            _trace_flow_id: TraceFlowId,
1757        ) -> Result<(), zx::Status> {
1758            if opts.flags.contains(WriteFlags::PRE_BARRIER)
1759                && let Some(barrier_hook) = &self.barrier_hook
1760            {
1761                barrier_hook()?;
1762            }
1763            if let Some(write_hook) = &self.write_hook {
1764                write_hook(device_block_offset).await
1765            } else {
1766                unimplemented!();
1767            }
1768        }
1769
1770        async fn flush(&self, _trace_flow_id: TraceFlowId) -> Result<(), zx::Status> {
1771            Ok(())
1772        }
1773
1774        async fn trim(
1775            &self,
1776            _device_block_offset: u64,
1777            _block_count: u32,
1778            _trace_flow_id: TraceFlowId,
1779        ) -> Result<(), zx::Status> {
1780            unreachable!();
1781        }
1782
1783        async fn get_volume_info(
1784            &self,
1785        ) -> Result<(fblock::VolumeManagerInfo, fblock::VolumeInfo), zx::Status> {
1786            // Hang forever for the test_requests_dont_block_sessions test.
1787            let () = std::future::pending().await;
1788            unreachable!();
1789        }
1790    }
1791
1792    const BLOCK_SIZE: u32 = 512;
1793    const MAX_TRANSFER_BLOCKS: u32 = 10;
1794
1795    fn test_device_info() -> DeviceInfo {
1796        DeviceInfo::Partition(PartitionInfo {
1797            device_flags: fblock::DeviceFlag::READONLY
1798                | fblock::DeviceFlag::BARRIER_SUPPORT
1799                | fblock::DeviceFlag::FUA_SUPPORT,
1800            max_transfer_blocks: NonZero::new(MAX_TRANSFER_BLOCKS),
1801            start_block_offset: Some(0),
1802            block_count: 100,
1803            type_guid: [1; 16],
1804            instance_guid: [2; 16],
1805            name: "foo".to_string(),
1806            flags: Some(0xabcd),
1807        })
1808    }
1809
1810    #[fuchsia::test]
1811    async fn test_barriers_ordering() {
1812        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
1813        let vmo = zx::Vmo::create(2 * zx::system_get_page_size() as u64).unwrap();
1814        let barrier_called = Arc::new(AtomicBool::new(false));
1815
1816        futures::join!(
1817            async move {
1818                let barrier_called_clone = barrier_called.clone();
1819                let block_server = BlockServer::new(
1820                    BLOCK_SIZE,
1821                    Arc::new(MockInterface {
1822                        barrier_hook: Some(Box::new(move || {
1823                            barrier_called.store(true, Ordering::Relaxed);
1824                            Ok(())
1825                        })),
1826                        write_hook: Some(Box::new(move |device_block_offset| {
1827                            let barrier_called = barrier_called_clone.clone();
1828                            Box::pin(async move {
1829                                // The sleep allows the server to reorder the fifo requests.
1830                                if device_block_offset % 2 == 0 {
1831                                    fasync::Timer::new(fasync::MonotonicInstant::after(
1832                                        zx::MonotonicDuration::from_millis(200),
1833                                    ))
1834                                    .await;
1835                                }
1836                                assert!(barrier_called.load(Ordering::Relaxed));
1837                                Ok(())
1838                            })
1839                        })),
1840                        ..MockInterface::default()
1841                    }),
1842                );
1843                block_server.handle_requests(stream).await.unwrap();
1844            },
1845            async move {
1846                let (session_proxy, server) = fidl::endpoints::create_proxy();
1847
1848                proxy.open_session(server).unwrap();
1849
1850                let vmo_id = session_proxy
1851                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
1852                    .await
1853                    .unwrap()
1854                    .unwrap();
1855                assert_ne!(vmo_id.id, 0);
1856
1857                let mut fifo =
1858                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
1859                let (mut reader, mut writer) = fifo.async_io();
1860
1861                writer
1862                    .write_entries(&BlockFifoRequest {
1863                        command: BlockFifoCommand {
1864                            opcode: BlockOpcode::Write.into_primitive(),
1865                            flags: BlockIoFlag::PRE_BARRIER.bits(),
1866                            ..Default::default()
1867                        },
1868                        vmoid: vmo_id.id,
1869                        dev_offset: 0,
1870                        length: 5,
1871                        vmo_offset: 6,
1872                        ..Default::default()
1873                    })
1874                    .await
1875                    .unwrap();
1876
1877                for i in 0..10 {
1878                    writer
1879                        .write_entries(&BlockFifoRequest {
1880                            command: BlockFifoCommand {
1881                                opcode: BlockOpcode::Write.into_primitive(),
1882                                ..Default::default()
1883                            },
1884                            vmoid: vmo_id.id,
1885                            dev_offset: i + 1,
1886                            length: 5,
1887                            vmo_offset: 6,
1888                            ..Default::default()
1889                        })
1890                        .await
1891                        .unwrap();
1892                }
1893                for _ in 0..11 {
1894                    let mut response = BlockFifoResponse::default();
1895                    reader.read_entries(&mut response).await.unwrap();
1896                    assert_eq!(response.status, zx::sys::ZX_OK);
1897                }
1898
1899                std::mem::drop(proxy);
1900            }
1901        );
1902    }
1903
1904    #[fuchsia::test]
1905    async fn test_info() {
1906        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
1907
1908        futures::join!(
1909            async {
1910                let block_server = BlockServer::new(BLOCK_SIZE, Arc::new(MockInterface::default()));
1911                block_server.handle_requests(stream).await.unwrap();
1912            },
1913            async {
1914                let expected_info = test_device_info();
1915                let partition_info = if let DeviceInfo::Partition(info) = &expected_info {
1916                    info
1917                } else {
1918                    unreachable!()
1919                };
1920
1921                let block_info = proxy.get_info().await.unwrap().unwrap();
1922                assert_eq!(block_info.block_count, expected_info.block_count().unwrap());
1923                assert_eq!(
1924                    block_info.flags,
1925                    fblock::DeviceFlag::READONLY
1926                        | fblock::DeviceFlag::ZSTD_DECOMPRESSION_SUPPORT
1927                        | fblock::DeviceFlag::BARRIER_SUPPORT
1928                        | fblock::DeviceFlag::FUA_SUPPORT
1929                );
1930
1931                assert_eq!(block_info.max_transfer_size, MAX_TRANSFER_BLOCKS * BLOCK_SIZE);
1932
1933                let (status, type_guid) = proxy.get_type_guid().await.unwrap();
1934                assert_eq!(status, zx::sys::ZX_OK);
1935                assert_eq!(&type_guid.as_ref().unwrap().value, &partition_info.type_guid);
1936
1937                let (status, instance_guid) = proxy.get_instance_guid().await.unwrap();
1938                assert_eq!(status, zx::sys::ZX_OK);
1939                assert_eq!(&instance_guid.as_ref().unwrap().value, &partition_info.instance_guid);
1940
1941                let (status, name) = proxy.get_name().await.unwrap();
1942                assert_eq!(status, zx::sys::ZX_OK);
1943                assert_eq!(name.as_ref(), Some(&partition_info.name));
1944
1945                let metadata = proxy.get_metadata().await.unwrap().expect("get_flags failed");
1946                assert_eq!(metadata.name, name);
1947                assert_eq!(metadata.type_guid.as_ref(), type_guid.as_deref());
1948                assert_eq!(metadata.instance_guid.as_ref(), instance_guid.as_deref());
1949                let expected_start = partition_info.start_block_offset.or(Some(0));
1950                assert_eq!(metadata.start_block_offset, expected_start);
1951                assert_eq!(metadata.num_blocks, Some(partition_info.block_count));
1952                assert_eq!(metadata.flags, partition_info.flags);
1953
1954                std::mem::drop(proxy);
1955            }
1956        );
1957    }
1958
1959    #[fuchsia::test]
1960    async fn test_attach_vmo() {
1961        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
1962
1963        let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
1964        let koid = vmo.koid().unwrap();
1965
1966        futures::join!(
1967            async {
1968                let block_server = BlockServer::new(
1969                    BLOCK_SIZE,
1970                    Arc::new(MockInterface {
1971                        read_hook: Some(Box::new(move |_, _, vmo, _| {
1972                            assert_eq!(vmo.koid().unwrap(), koid);
1973                            Box::pin(async { Ok(()) })
1974                        })),
1975                        ..MockInterface::default()
1976                    }),
1977                );
1978                block_server.handle_requests(stream).await.unwrap();
1979            },
1980            async move {
1981                let (session_proxy, server) = fidl::endpoints::create_proxy();
1982
1983                proxy.open_session(server).unwrap();
1984
1985                let vmo_id = session_proxy
1986                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
1987                    .await
1988                    .unwrap()
1989                    .unwrap();
1990                assert_ne!(vmo_id.id, 0);
1991
1992                let mut fifo =
1993                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
1994                let (mut reader, mut writer) = fifo.async_io();
1995
1996                // Keep attaching VMOs until we eventually hit the maximum.
1997                let mut count = 1;
1998                loop {
1999                    match session_proxy
2000                        .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2001                        .await
2002                        .unwrap()
2003                    {
2004                        Ok(vmo_id) => assert_ne!(vmo_id.id, 0),
2005                        Err(e) => {
2006                            assert_eq!(e, zx::sys::ZX_ERR_NO_RESOURCES);
2007                            break;
2008                        }
2009                    }
2010
2011                    // Only test every 10 to keep test time down.
2012                    if count % 10 == 0 {
2013                        writer
2014                            .write_entries(&BlockFifoRequest {
2015                                command: BlockFifoCommand {
2016                                    opcode: BlockOpcode::Read.into_primitive(),
2017                                    ..Default::default()
2018                                },
2019                                vmoid: vmo_id.id,
2020                                length: 1,
2021                                ..Default::default()
2022                            })
2023                            .await
2024                            .unwrap();
2025
2026                        let mut response = BlockFifoResponse::default();
2027                        reader.read_entries(&mut response).await.unwrap();
2028                        assert_eq!(response.status, zx::sys::ZX_OK);
2029                    }
2030
2031                    count += 1;
2032                }
2033
2034                assert_eq!(count, u16::MAX as u64);
2035
2036                // Detach the original VMO, and make sure we can then attach another one.
2037                writer
2038                    .write_entries(&BlockFifoRequest {
2039                        command: BlockFifoCommand {
2040                            opcode: BlockOpcode::CloseVmo.into_primitive(),
2041                            ..Default::default()
2042                        },
2043                        vmoid: vmo_id.id,
2044                        ..Default::default()
2045                    })
2046                    .await
2047                    .unwrap();
2048
2049                let mut response = BlockFifoResponse::default();
2050                reader.read_entries(&mut response).await.unwrap();
2051                assert_eq!(response.status, zx::sys::ZX_OK);
2052
2053                let new_vmo_id = session_proxy
2054                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2055                    .await
2056                    .unwrap()
2057                    .unwrap();
2058                // It should reuse the same ID.
2059                assert_eq!(new_vmo_id.id, vmo_id.id);
2060
2061                std::mem::drop(proxy);
2062            }
2063        );
2064    }
2065
2066    #[fuchsia::test]
2067    async fn test_attach_resizable_vmo_fails() {
2068        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2069        let resizable_vmo = zx::Vmo::create_with_opts(zx::VmoOptions::RESIZABLE, 4096).unwrap();
2070
2071        futures::join!(
2072            async {
2073                let block_server = BlockServer::new(BLOCK_SIZE, Arc::new(MockInterface::default()));
2074                block_server.handle_requests(stream).await.unwrap();
2075            },
2076            async move {
2077                let (session_proxy, server) = fidl::endpoints::create_proxy();
2078                proxy.open_session(server).unwrap();
2079                let res = session_proxy.attach_vmo(resizable_vmo).await.unwrap();
2080                assert_eq!(res, Err(zx::sys::ZX_ERR_INVALID_ARGS));
2081                std::mem::drop(proxy);
2082            }
2083        );
2084    }
2085
2086    #[fuchsia::test]
2087    async fn test_close() {
2088        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2089
2090        let mut server = std::pin::pin!(
2091            async {
2092                let block_server = BlockServer::new(BLOCK_SIZE, Arc::new(MockInterface::default()));
2093                block_server.handle_requests(stream).await.unwrap();
2094            }
2095            .fuse()
2096        );
2097
2098        let mut client = std::pin::pin!(
2099            async {
2100                let (session_proxy, server) = fidl::endpoints::create_proxy();
2101
2102                proxy.open_session(server).unwrap();
2103
2104                // Dropping the proxy should not cause the session to terminate because the session
2105                // is still live.
2106                std::mem::drop(proxy);
2107
2108                session_proxy.close().await.unwrap().unwrap();
2109
2110                // Keep the session alive.  Calling `close` should cause the server to terminate.
2111                let _: () = std::future::pending().await;
2112            }
2113            .fuse()
2114        );
2115
2116        futures::select!(
2117            _ = server => {}
2118            _ = client => unreachable!(),
2119        );
2120    }
2121
2122    #[derive(Default)]
2123    struct IoMockInterface {
2124        do_checks: bool,
2125        expected_op: Arc<Mutex<Option<ExpectedOp>>>,
2126        return_errors: bool,
2127    }
2128
2129    #[derive(Debug)]
2130    enum ExpectedOp {
2131        Read(u64, u32, u64),
2132        Write(u64, u32, u64),
2133        Trim(u64, u32),
2134        Flush,
2135    }
2136
2137    impl super::async_interface::Interface for IoMockInterface {
2138        async fn on_attach_vmo(&self, _vmo: &zx::Vmo) -> Result<(), zx::Status> {
2139            Ok(())
2140        }
2141
2142        fn get_info(&self) -> Cow<'_, DeviceInfo> {
2143            Cow::Owned(DeviceInfo::Block(crate::BlockInfo {
2144                block_count: 100,
2145                ..Default::default()
2146            }))
2147        }
2148
2149        async fn read(
2150            &self,
2151            device_block_offset: u64,
2152            block_count: u32,
2153            _vmo: &Arc<zx::Vmo>,
2154            vmo_offset: u64,
2155            _opts: ReadOptions,
2156            _trace_flow_id: TraceFlowId,
2157        ) -> Result<(), zx::Status> {
2158            if self.return_errors {
2159                Err(zx::Status::INTERNAL)
2160            } else {
2161                if self.do_checks {
2162                    assert_matches!(
2163                        self.expected_op.lock().take(),
2164                        Some(ExpectedOp::Read(a, b, c)) if device_block_offset == a &&
2165                            block_count == b && vmo_offset / BLOCK_SIZE as u64 == c,
2166                        "Read {device_block_offset} {block_count} {vmo_offset}"
2167                    );
2168                }
2169                Ok(())
2170            }
2171        }
2172
2173        async fn write(
2174            &self,
2175            device_block_offset: u64,
2176            block_count: u32,
2177            _vmo: &Arc<zx::Vmo>,
2178            vmo_offset: u64,
2179            _write_opts: WriteOptions,
2180            _trace_flow_id: TraceFlowId,
2181        ) -> Result<(), zx::Status> {
2182            if self.return_errors {
2183                Err(zx::Status::NOT_SUPPORTED)
2184            } else {
2185                if self.do_checks {
2186                    assert_matches!(
2187                        self.expected_op.lock().take(),
2188                        Some(ExpectedOp::Write(a, b, c)) if device_block_offset == a &&
2189                            block_count == b && vmo_offset / BLOCK_SIZE as u64 == c,
2190                        "Write {device_block_offset} {block_count} {vmo_offset}"
2191                    );
2192                }
2193                Ok(())
2194            }
2195        }
2196
2197        async fn flush(&self, _trace_flow_id: TraceFlowId) -> Result<(), zx::Status> {
2198            if self.return_errors {
2199                Err(zx::Status::NO_RESOURCES)
2200            } else {
2201                if self.do_checks {
2202                    assert_matches!(self.expected_op.lock().take(), Some(ExpectedOp::Flush));
2203                }
2204                Ok(())
2205            }
2206        }
2207
2208        async fn trim(
2209            &self,
2210            device_block_offset: u64,
2211            block_count: u32,
2212            _trace_flow_id: TraceFlowId,
2213        ) -> Result<(), zx::Status> {
2214            if self.return_errors {
2215                Err(zx::Status::NO_MEMORY)
2216            } else {
2217                if self.do_checks {
2218                    assert_matches!(
2219                        self.expected_op.lock().take(),
2220                        Some(ExpectedOp::Trim(a, b)) if device_block_offset == a &&
2221                            block_count == b,
2222                        "Trim {device_block_offset} {block_count}"
2223                    );
2224                }
2225                Ok(())
2226            }
2227        }
2228    }
2229
2230    #[fuchsia::test]
2231    async fn test_io() {
2232        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2233
2234        let expected_op = Arc::new(Mutex::new(None));
2235        let expected_op_clone = expected_op.clone();
2236
2237        let server = async {
2238            let block_server = BlockServer::new(
2239                BLOCK_SIZE,
2240                Arc::new(IoMockInterface {
2241                    return_errors: false,
2242                    do_checks: true,
2243                    expected_op: expected_op_clone,
2244                }),
2245            );
2246            block_server.handle_requests(stream).await.unwrap();
2247        };
2248
2249        let client = async move {
2250            let (session_proxy, server) = fidl::endpoints::create_proxy();
2251
2252            proxy.open_session(server).unwrap();
2253
2254            let vmo = zx::Vmo::create(2 * zx::system_get_page_size() as u64).unwrap();
2255            let vmo_id = session_proxy
2256                .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2257                .await
2258                .unwrap()
2259                .unwrap();
2260
2261            let mut fifo =
2262                fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
2263            let (mut reader, mut writer) = fifo.async_io();
2264
2265            // READ
2266            *expected_op.lock() = Some(ExpectedOp::Read(1, 2, 3));
2267            writer
2268                .write_entries(&BlockFifoRequest {
2269                    command: BlockFifoCommand {
2270                        opcode: BlockOpcode::Read.into_primitive(),
2271                        ..Default::default()
2272                    },
2273                    vmoid: vmo_id.id,
2274                    dev_offset: 1,
2275                    length: 2,
2276                    vmo_offset: 3,
2277                    ..Default::default()
2278                })
2279                .await
2280                .unwrap();
2281
2282            let mut response = BlockFifoResponse::default();
2283            reader.read_entries(&mut response).await.unwrap();
2284            assert_eq!(response.status, zx::sys::ZX_OK);
2285
2286            // WRITE
2287            *expected_op.lock() = Some(ExpectedOp::Write(4, 5, 6));
2288            writer
2289                .write_entries(&BlockFifoRequest {
2290                    command: BlockFifoCommand {
2291                        opcode: BlockOpcode::Write.into_primitive(),
2292                        ..Default::default()
2293                    },
2294                    vmoid: vmo_id.id,
2295                    dev_offset: 4,
2296                    length: 5,
2297                    vmo_offset: 6,
2298                    ..Default::default()
2299                })
2300                .await
2301                .unwrap();
2302
2303            let mut response = BlockFifoResponse::default();
2304            reader.read_entries(&mut response).await.unwrap();
2305            assert_eq!(response.status, zx::sys::ZX_OK);
2306
2307            // FLUSH
2308            *expected_op.lock() = Some(ExpectedOp::Flush);
2309            writer
2310                .write_entries(&BlockFifoRequest {
2311                    command: BlockFifoCommand {
2312                        opcode: BlockOpcode::Flush.into_primitive(),
2313                        ..Default::default()
2314                    },
2315                    ..Default::default()
2316                })
2317                .await
2318                .unwrap();
2319
2320            reader.read_entries(&mut response).await.unwrap();
2321            assert_eq!(response.status, zx::sys::ZX_OK);
2322
2323            // TRIM
2324            *expected_op.lock() = Some(ExpectedOp::Trim(7, 8));
2325            writer
2326                .write_entries(&BlockFifoRequest {
2327                    command: BlockFifoCommand {
2328                        opcode: BlockOpcode::Trim.into_primitive(),
2329                        ..Default::default()
2330                    },
2331                    dev_offset: 7,
2332                    length: 8,
2333                    ..Default::default()
2334                })
2335                .await
2336                .unwrap();
2337
2338            reader.read_entries(&mut response).await.unwrap();
2339            assert_eq!(response.status, zx::sys::ZX_OK);
2340
2341            std::mem::drop(proxy);
2342        };
2343
2344        futures::join!(server, client);
2345    }
2346
2347    #[fuchsia::test]
2348    async fn test_io_errors() {
2349        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2350
2351        futures::join!(
2352            async {
2353                let block_server = BlockServer::new(
2354                    BLOCK_SIZE,
2355                    Arc::new(IoMockInterface {
2356                        return_errors: true,
2357                        do_checks: false,
2358                        expected_op: Arc::new(Mutex::new(None)),
2359                    }),
2360                );
2361                block_server.handle_requests(stream).await.unwrap();
2362            },
2363            async move {
2364                let (session_proxy, server) = fidl::endpoints::create_proxy();
2365
2366                proxy.open_session(server).unwrap();
2367
2368                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
2369                let vmo_id = session_proxy
2370                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2371                    .await
2372                    .unwrap()
2373                    .unwrap();
2374
2375                let mut fifo =
2376                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
2377                let (mut reader, mut writer) = fifo.async_io();
2378
2379                // READ
2380                writer
2381                    .write_entries(&BlockFifoRequest {
2382                        command: BlockFifoCommand {
2383                            opcode: BlockOpcode::Read.into_primitive(),
2384                            ..Default::default()
2385                        },
2386                        vmoid: vmo_id.id,
2387                        length: 1,
2388                        reqid: 1,
2389                        ..Default::default()
2390                    })
2391                    .await
2392                    .unwrap();
2393
2394                let mut response = BlockFifoResponse::default();
2395                reader.read_entries(&mut response).await.unwrap();
2396                assert_eq!(response.status, zx::sys::ZX_ERR_INTERNAL);
2397
2398                // WRITE
2399                writer
2400                    .write_entries(&BlockFifoRequest {
2401                        command: BlockFifoCommand {
2402                            opcode: BlockOpcode::Write.into_primitive(),
2403                            ..Default::default()
2404                        },
2405                        vmoid: vmo_id.id,
2406                        length: 1,
2407                        reqid: 2,
2408                        ..Default::default()
2409                    })
2410                    .await
2411                    .unwrap();
2412
2413                reader.read_entries(&mut response).await.unwrap();
2414                assert_eq!(response.status, zx::sys::ZX_ERR_NOT_SUPPORTED);
2415
2416                // FLUSH
2417                writer
2418                    .write_entries(&BlockFifoRequest {
2419                        command: BlockFifoCommand {
2420                            opcode: BlockOpcode::Flush.into_primitive(),
2421                            ..Default::default()
2422                        },
2423                        reqid: 3,
2424                        ..Default::default()
2425                    })
2426                    .await
2427                    .unwrap();
2428
2429                reader.read_entries(&mut response).await.unwrap();
2430                assert_eq!(response.status, zx::sys::ZX_ERR_NO_RESOURCES);
2431
2432                // TRIM
2433                writer
2434                    .write_entries(&BlockFifoRequest {
2435                        command: BlockFifoCommand {
2436                            opcode: BlockOpcode::Trim.into_primitive(),
2437                            ..Default::default()
2438                        },
2439                        reqid: 4,
2440                        length: 1,
2441                        ..Default::default()
2442                    })
2443                    .await
2444                    .unwrap();
2445
2446                reader.read_entries(&mut response).await.unwrap();
2447                assert_eq!(response.status, zx::sys::ZX_ERR_NO_MEMORY);
2448
2449                std::mem::drop(proxy);
2450            }
2451        );
2452    }
2453
2454    #[fuchsia::test]
2455    async fn test_invalid_args() {
2456        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2457
2458        futures::join!(
2459            async {
2460                let block_server = BlockServer::new(
2461                    BLOCK_SIZE,
2462                    Arc::new(IoMockInterface {
2463                        return_errors: false,
2464                        do_checks: false,
2465                        expected_op: Arc::new(Mutex::new(None)),
2466                    }),
2467                );
2468                block_server.handle_requests(stream).await.unwrap();
2469            },
2470            async move {
2471                let (session_proxy, server) = fidl::endpoints::create_proxy();
2472
2473                proxy.open_session(server).unwrap();
2474
2475                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
2476                let vmo_id = session_proxy
2477                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2478                    .await
2479                    .unwrap()
2480                    .unwrap();
2481
2482                let mut fifo =
2483                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
2484
2485                async fn test(
2486                    fifo: &mut fasync::Fifo<BlockFifoResponse, BlockFifoRequest>,
2487                    request: BlockFifoRequest,
2488                ) -> Result<(), zx::Status> {
2489                    let (mut reader, mut writer) = fifo.async_io();
2490                    writer.write_entries(&request).await.unwrap();
2491                    let mut response = BlockFifoResponse::default();
2492                    reader.read_entries(&mut response).await.unwrap();
2493                    zx::Status::ok(response.status)
2494                }
2495
2496                // READ
2497
2498                let good_read_request = || BlockFifoRequest {
2499                    command: BlockFifoCommand {
2500                        opcode: BlockOpcode::Read.into_primitive(),
2501                        ..Default::default()
2502                    },
2503                    length: 1,
2504                    vmoid: vmo_id.id,
2505                    ..Default::default()
2506                };
2507
2508                assert_eq!(
2509                    test(
2510                        &mut fifo,
2511                        BlockFifoRequest { vmoid: vmo_id.id + 1, ..good_read_request() }
2512                    )
2513                    .await,
2514                    Err(zx::Status::IO)
2515                );
2516
2517                assert_eq!(
2518                    test(
2519                        &mut fifo,
2520                        BlockFifoRequest {
2521                            vmo_offset: 0xffff_ffff_ffff_ffff,
2522                            ..good_read_request()
2523                        }
2524                    )
2525                    .await,
2526                    Err(zx::Status::OUT_OF_RANGE)
2527                );
2528
2529                assert_eq!(
2530                    test(
2531                        &mut fifo,
2532                        BlockFifoRequest {
2533                            vmo_offset: 0x007f_ffff_ffff_ffff,
2534                            length: 2,
2535                            ..good_read_request()
2536                        }
2537                    )
2538                    .await,
2539                    Err(zx::Status::OUT_OF_RANGE)
2540                );
2541
2542                assert_eq!(
2543                    test(&mut fifo, BlockFifoRequest { length: 0, ..good_read_request() }).await,
2544                    Err(zx::Status::INVALID_ARGS)
2545                );
2546
2547                assert_eq!(
2548                    test(
2549                        &mut fifo,
2550                        BlockFifoRequest { vmo_offset: 8, length: 1, ..good_read_request() }
2551                    )
2552                    .await,
2553                    Err(zx::Status::OUT_OF_RANGE)
2554                );
2555
2556                // WRITE
2557
2558                let good_write_request = || BlockFifoRequest {
2559                    command: BlockFifoCommand {
2560                        opcode: BlockOpcode::Write.into_primitive(),
2561                        ..Default::default()
2562                    },
2563                    length: 1,
2564                    vmoid: vmo_id.id,
2565                    ..Default::default()
2566                };
2567
2568                assert_eq!(
2569                    test(
2570                        &mut fifo,
2571                        BlockFifoRequest { vmoid: vmo_id.id + 1, ..good_write_request() }
2572                    )
2573                    .await,
2574                    Err(zx::Status::IO)
2575                );
2576
2577                assert_eq!(
2578                    test(
2579                        &mut fifo,
2580                        BlockFifoRequest {
2581                            vmo_offset: 0xffff_ffff_ffff_ffff,
2582                            ..good_write_request()
2583                        }
2584                    )
2585                    .await,
2586                    Err(zx::Status::OUT_OF_RANGE)
2587                );
2588
2589                assert_eq!(
2590                    test(
2591                        &mut fifo,
2592                        BlockFifoRequest {
2593                            vmo_offset: 0x007f_ffff_ffff_ffff,
2594                            length: 2,
2595                            ..good_write_request()
2596                        }
2597                    )
2598                    .await,
2599                    Err(zx::Status::OUT_OF_RANGE)
2600                );
2601
2602                assert_eq!(
2603                    test(&mut fifo, BlockFifoRequest { length: 0, ..good_write_request() }).await,
2604                    Err(zx::Status::INVALID_ARGS)
2605                );
2606
2607                assert_eq!(
2608                    test(
2609                        &mut fifo,
2610                        BlockFifoRequest { vmo_offset: 8, length: 1, ..good_write_request() }
2611                    )
2612                    .await,
2613                    Err(zx::Status::OUT_OF_RANGE)
2614                );
2615
2616                // CLOSE VMO
2617
2618                assert_eq!(
2619                    test(
2620                        &mut fifo,
2621                        BlockFifoRequest {
2622                            command: BlockFifoCommand {
2623                                opcode: BlockOpcode::CloseVmo.into_primitive(),
2624                                ..Default::default()
2625                            },
2626                            vmoid: vmo_id.id + 1,
2627                            ..Default::default()
2628                        }
2629                    )
2630                    .await,
2631                    Err(zx::Status::IO)
2632                );
2633
2634                std::mem::drop(proxy);
2635            }
2636        );
2637    }
2638
2639    #[fuchsia::test]
2640    async fn test_concurrent_requests() {
2641        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2642
2643        let waiting_readers = Arc::new(Mutex::new(Vec::new()));
2644        let waiting_readers_clone = waiting_readers.clone();
2645
2646        futures::join!(
2647            async move {
2648                let block_server = BlockServer::new(
2649                    BLOCK_SIZE,
2650                    Arc::new(MockInterface {
2651                        read_hook: Some(Box::new(move |dev_block_offset, _, _, _| {
2652                            let (tx, rx) = oneshot::channel();
2653                            waiting_readers_clone.lock().push((dev_block_offset as u32, tx));
2654                            Box::pin(async move {
2655                                let _ = rx.await;
2656                                Ok(())
2657                            })
2658                        })),
2659                        ..MockInterface::default()
2660                    }),
2661                );
2662                block_server.handle_requests(stream).await.unwrap();
2663            },
2664            async move {
2665                let (session_proxy, server) = fidl::endpoints::create_proxy();
2666
2667                proxy.open_session(server).unwrap();
2668
2669                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
2670                let vmo_id = session_proxy
2671                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2672                    .await
2673                    .unwrap()
2674                    .unwrap();
2675
2676                let mut fifo =
2677                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
2678                let (mut reader, mut writer) = fifo.async_io();
2679
2680                writer
2681                    .write_entries(&BlockFifoRequest {
2682                        command: BlockFifoCommand {
2683                            opcode: BlockOpcode::Read.into_primitive(),
2684                            ..Default::default()
2685                        },
2686                        reqid: 1,
2687                        dev_offset: 1, // Intentionally use the same as `reqid`.
2688                        vmoid: vmo_id.id,
2689                        length: 1,
2690                        ..Default::default()
2691                    })
2692                    .await
2693                    .unwrap();
2694
2695                writer
2696                    .write_entries(&BlockFifoRequest {
2697                        command: BlockFifoCommand {
2698                            opcode: BlockOpcode::Read.into_primitive(),
2699                            ..Default::default()
2700                        },
2701                        reqid: 2,
2702                        dev_offset: 2,
2703                        vmoid: vmo_id.id,
2704                        length: 1,
2705                        ..Default::default()
2706                    })
2707                    .await
2708                    .unwrap();
2709
2710                // Wait till both those entries are pending.
2711                poll_fn(|cx: &mut Context<'_>| {
2712                    if waiting_readers.lock().len() == 2 {
2713                        Poll::Ready(())
2714                    } else {
2715                        // Yield to the executor.
2716                        cx.waker().wake_by_ref();
2717                        Poll::Pending
2718                    }
2719                })
2720                .await;
2721
2722                let mut response = BlockFifoResponse::default();
2723                assert!(futures::poll!(pin!(reader.read_entries(&mut response))).is_pending());
2724
2725                let (id, tx) = waiting_readers.lock().pop().unwrap();
2726                tx.send(()).unwrap();
2727
2728                reader.read_entries(&mut response).await.unwrap();
2729                assert_eq!(response.status, zx::sys::ZX_OK);
2730                assert_eq!(response.reqid, id);
2731
2732                assert!(futures::poll!(pin!(reader.read_entries(&mut response))).is_pending());
2733
2734                let (id, tx) = waiting_readers.lock().pop().unwrap();
2735                tx.send(()).unwrap();
2736
2737                reader.read_entries(&mut response).await.unwrap();
2738                assert_eq!(response.status, zx::sys::ZX_OK);
2739                assert_eq!(response.reqid, id);
2740            }
2741        );
2742    }
2743
2744    #[fuchsia::test]
2745    async fn test_session_close_is_synchronous() {
2746        use futures::{FutureExt as _, StreamExt as _};
2747
2748        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2749
2750        let (start_tx, mut start_rx) = futures::channel::mpsc::channel(1);
2751        let (finish_tx, finish_rx) = futures::channel::oneshot::channel();
2752        let finish_rx = Arc::new(Mutex::new(Some(finish_rx)));
2753
2754        futures::join!(
2755            async move {
2756                let block_server = BlockServer::new(
2757                    BLOCK_SIZE,
2758                    Arc::new(MockInterface {
2759                        read_hook: Some(Box::new(move |_, _, _, _| {
2760                            let mut start_tx = start_tx.clone();
2761                            let finish_rx = finish_rx.lock().take().unwrap();
2762                            Box::pin(async move {
2763                                start_tx.try_send(()).unwrap();
2764                                let _ = finish_rx.await;
2765                                Ok(())
2766                            })
2767                        })),
2768                        ..MockInterface::default()
2769                    }),
2770                );
2771                block_server.handle_requests(stream).await.unwrap();
2772            },
2773            async move {
2774                let (session_proxy, server) = fidl::endpoints::create_proxy();
2775                proxy.open_session(server).unwrap();
2776
2777                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
2778                let vmo_id = session_proxy
2779                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2780                    .await
2781                    .unwrap()
2782                    .unwrap();
2783
2784                let mut fifo = fasync::Fifo::<BlockFifoResponse, BlockFifoRequest>::from_fifo(
2785                    session_proxy.get_fifo().await.unwrap().unwrap(),
2786                );
2787                let (_reader, mut writer) = fifo.async_io();
2788
2789                writer
2790                    .write_entries(&BlockFifoRequest {
2791                        command: BlockFifoCommand {
2792                            opcode: BlockOpcode::Read.into_primitive(),
2793                            ..Default::default()
2794                        },
2795                        reqid: 1,
2796                        vmoid: vmo_id.id,
2797                        length: 1,
2798                        ..Default::default()
2799                    })
2800                    .await
2801                    .unwrap();
2802
2803                // Wait for the read to actually start.
2804                start_rx.next().await.unwrap();
2805
2806                // The close request shouldn't complete yet because the read is still hanging.
2807                let mut close_fut = std::pin::pin!(session_proxy.close().fuse());
2808                let mut timer_fut = std::pin::pin!(
2809                    fasync::Timer::new(std::time::Duration::from_millis(100)).fuse()
2810                );
2811                futures::select! {
2812                    res = close_fut => panic!("close completed too early: {:?}", res),
2813                    _ = timer_fut => {}
2814                }
2815
2816                // Finish the pending request.
2817                finish_tx.send(()).unwrap();
2818
2819                // Verify that close() now completes.
2820                close_fut.await.unwrap().unwrap();
2821
2822                std::mem::drop(proxy);
2823            }
2824        );
2825    }
2826
2827    #[fuchsia::test]
2828    async fn test_groups() {
2829        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2830
2831        futures::join!(
2832            async move {
2833                let block_server = BlockServer::new(
2834                    BLOCK_SIZE,
2835                    Arc::new(MockInterface {
2836                        read_hook: Some(Box::new(move |_, _, _, _| Box::pin(async { Ok(()) }))),
2837                        ..MockInterface::default()
2838                    }),
2839                );
2840                block_server.handle_requests(stream).await.unwrap();
2841            },
2842            async move {
2843                let (session_proxy, server) = fidl::endpoints::create_proxy();
2844
2845                proxy.open_session(server).unwrap();
2846
2847                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
2848                let vmo_id = session_proxy
2849                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2850                    .await
2851                    .unwrap()
2852                    .unwrap();
2853
2854                let mut fifo =
2855                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
2856                let (mut reader, mut writer) = fifo.async_io();
2857
2858                writer
2859                    .write_entries(&BlockFifoRequest {
2860                        command: BlockFifoCommand {
2861                            opcode: BlockOpcode::Read.into_primitive(),
2862                            flags: BlockIoFlag::GROUP_ITEM.bits(),
2863                            ..Default::default()
2864                        },
2865                        group: 1,
2866                        vmoid: vmo_id.id,
2867                        length: 1,
2868                        ..Default::default()
2869                    })
2870                    .await
2871                    .unwrap();
2872
2873                writer
2874                    .write_entries(&BlockFifoRequest {
2875                        command: BlockFifoCommand {
2876                            opcode: BlockOpcode::Read.into_primitive(),
2877                            flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
2878                            ..Default::default()
2879                        },
2880                        reqid: 2,
2881                        group: 1,
2882                        vmoid: vmo_id.id,
2883                        length: 1,
2884                        ..Default::default()
2885                    })
2886                    .await
2887                    .unwrap();
2888
2889                let mut response = BlockFifoResponse::default();
2890                reader.read_entries(&mut response).await.unwrap();
2891                assert_eq!(response.status, zx::sys::ZX_OK);
2892                assert_eq!(response.reqid, 2);
2893                assert_eq!(response.group, 1);
2894            }
2895        );
2896    }
2897
2898    #[fuchsia::test]
2899    async fn test_group_error() {
2900        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
2901
2902        let counter = Arc::new(AtomicU64::new(0));
2903        let counter_clone = counter.clone();
2904
2905        futures::join!(
2906            async move {
2907                let block_server = BlockServer::new(
2908                    BLOCK_SIZE,
2909                    Arc::new(MockInterface {
2910                        read_hook: Some(Box::new(move |_, _, _, _| {
2911                            counter_clone.fetch_add(1, Ordering::Relaxed);
2912                            Box::pin(async { Err(zx::Status::BAD_STATE) })
2913                        })),
2914                        ..MockInterface::default()
2915                    }),
2916                );
2917                block_server.handle_requests(stream).await.unwrap();
2918            },
2919            async move {
2920                let (session_proxy, server) = fidl::endpoints::create_proxy();
2921
2922                proxy.open_session(server).unwrap();
2923
2924                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
2925                let vmo_id = session_proxy
2926                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
2927                    .await
2928                    .unwrap()
2929                    .unwrap();
2930
2931                let mut fifo =
2932                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
2933                let (mut reader, mut writer) = fifo.async_io();
2934
2935                writer
2936                    .write_entries(&BlockFifoRequest {
2937                        command: BlockFifoCommand {
2938                            opcode: BlockOpcode::Read.into_primitive(),
2939                            flags: BlockIoFlag::GROUP_ITEM.bits(),
2940                            ..Default::default()
2941                        },
2942                        group: 1,
2943                        vmoid: vmo_id.id,
2944                        length: 1,
2945                        ..Default::default()
2946                    })
2947                    .await
2948                    .unwrap();
2949
2950                // Wait until processed.
2951                poll_fn(|cx: &mut Context<'_>| {
2952                    if counter.load(Ordering::Relaxed) == 1 {
2953                        Poll::Ready(())
2954                    } else {
2955                        // Yield to the executor.
2956                        cx.waker().wake_by_ref();
2957                        Poll::Pending
2958                    }
2959                })
2960                .await;
2961
2962                let mut response = BlockFifoResponse::default();
2963                assert!(futures::poll!(pin!(reader.read_entries(&mut response))).is_pending());
2964
2965                writer
2966                    .write_entries(&BlockFifoRequest {
2967                        command: BlockFifoCommand {
2968                            opcode: BlockOpcode::Read.into_primitive(),
2969                            flags: BlockIoFlag::GROUP_ITEM.bits(),
2970                            ..Default::default()
2971                        },
2972                        group: 1,
2973                        vmoid: vmo_id.id,
2974                        length: 1,
2975                        ..Default::default()
2976                    })
2977                    .await
2978                    .unwrap();
2979
2980                writer
2981                    .write_entries(&BlockFifoRequest {
2982                        command: BlockFifoCommand {
2983                            opcode: BlockOpcode::Read.into_primitive(),
2984                            flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
2985                            ..Default::default()
2986                        },
2987                        reqid: 2,
2988                        group: 1,
2989                        vmoid: vmo_id.id,
2990                        length: 1,
2991                        ..Default::default()
2992                    })
2993                    .await
2994                    .unwrap();
2995
2996                reader.read_entries(&mut response).await.unwrap();
2997                assert_eq!(response.status, zx::sys::ZX_ERR_BAD_STATE);
2998                assert_eq!(response.reqid, 2);
2999                assert_eq!(response.group, 1);
3000
3001                assert!(futures::poll!(pin!(reader.read_entries(&mut response))).is_pending());
3002
3003                // Only the first request should have been processed.
3004                assert_eq!(counter.load(Ordering::Relaxed), 1);
3005            }
3006        );
3007    }
3008
3009    #[fuchsia::test]
3010    async fn test_group_with_two_lasts() {
3011        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3012
3013        let (tx, rx) = oneshot::channel();
3014
3015        futures::join!(
3016            async move {
3017                let rx = Mutex::new(Some(rx));
3018                let block_server = BlockServer::new(
3019                    BLOCK_SIZE,
3020                    Arc::new(MockInterface {
3021                        read_hook: Some(Box::new(move |_, _, _, _| {
3022                            let rx = rx.lock().take().unwrap();
3023                            Box::pin(async {
3024                                let _ = rx.await;
3025                                Ok(())
3026                            })
3027                        })),
3028                        ..MockInterface::default()
3029                    }),
3030                );
3031                block_server.handle_requests(stream).await.unwrap();
3032            },
3033            async move {
3034                let (session_proxy, server) = fidl::endpoints::create_proxy();
3035
3036                proxy.open_session(server).unwrap();
3037
3038                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
3039                let vmo_id = session_proxy
3040                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
3041                    .await
3042                    .unwrap()
3043                    .unwrap();
3044
3045                let mut fifo =
3046                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
3047                let (mut reader, mut writer) = fifo.async_io();
3048
3049                writer
3050                    .write_entries(&BlockFifoRequest {
3051                        command: BlockFifoCommand {
3052                            opcode: BlockOpcode::Read.into_primitive(),
3053                            flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
3054                            ..Default::default()
3055                        },
3056                        reqid: 1,
3057                        group: 1,
3058                        vmoid: vmo_id.id,
3059                        length: 1,
3060                        ..Default::default()
3061                    })
3062                    .await
3063                    .unwrap();
3064
3065                writer
3066                    .write_entries(&BlockFifoRequest {
3067                        command: BlockFifoCommand {
3068                            opcode: BlockOpcode::Read.into_primitive(),
3069                            flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
3070                            ..Default::default()
3071                        },
3072                        reqid: 2,
3073                        group: 1,
3074                        vmoid: vmo_id.id,
3075                        length: 1,
3076                        ..Default::default()
3077                    })
3078                    .await
3079                    .unwrap();
3080
3081                // Send an independent request to flush through the fifo.
3082                writer
3083                    .write_entries(&BlockFifoRequest {
3084                        command: BlockFifoCommand {
3085                            opcode: BlockOpcode::CloseVmo.into_primitive(),
3086                            ..Default::default()
3087                        },
3088                        reqid: 3,
3089                        vmoid: vmo_id.id,
3090                        ..Default::default()
3091                    })
3092                    .await
3093                    .unwrap();
3094
3095                // It should succeed.
3096                let mut response = BlockFifoResponse::default();
3097                reader.read_entries(&mut response).await.unwrap();
3098                assert_eq!(response.status, zx::sys::ZX_OK);
3099                assert_eq!(response.reqid, 3);
3100
3101                // Now release the original request.
3102                tx.send(()).unwrap();
3103
3104                // The response should be for the first message tagged as last, and it should be
3105                // an error because we sent two messages with the LAST marker.
3106                let mut response = BlockFifoResponse::default();
3107                reader.read_entries(&mut response).await.unwrap();
3108                assert_eq!(response.status, zx::sys::ZX_ERR_INVALID_ARGS);
3109                assert_eq!(response.reqid, 1);
3110                assert_eq!(response.group, 1);
3111            }
3112        );
3113    }
3114
3115    #[fuchsia::test(allow_stalls = false)]
3116    async fn test_requests_dont_block_sessions() {
3117        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3118
3119        let (tx, rx) = oneshot::channel();
3120
3121        fasync::Task::local(async move {
3122            let rx = Mutex::new(Some(rx));
3123            let block_server = BlockServer::new(
3124                BLOCK_SIZE,
3125                Arc::new(MockInterface {
3126                    read_hook: Some(Box::new(move |_, _, _, _| {
3127                        let rx = rx.lock().take().unwrap();
3128                        Box::pin(async {
3129                            let _ = rx.await;
3130                            Ok(())
3131                        })
3132                    })),
3133                    ..MockInterface::default()
3134                }),
3135            );
3136            block_server.handle_requests(stream).await.unwrap();
3137        })
3138        .detach();
3139
3140        let mut fut = pin!(async {
3141            let (session_proxy, server) = fidl::endpoints::create_proxy();
3142
3143            proxy.open_session(server).unwrap();
3144
3145            let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
3146            let vmo_id = session_proxy
3147                .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
3148                .await
3149                .unwrap()
3150                .unwrap();
3151
3152            let mut fifo =
3153                fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
3154            let (mut reader, mut writer) = fifo.async_io();
3155
3156            writer
3157                .write_entries(&BlockFifoRequest {
3158                    command: BlockFifoCommand {
3159                        opcode: BlockOpcode::Read.into_primitive(),
3160                        flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
3161                        ..Default::default()
3162                    },
3163                    reqid: 1,
3164                    group: 1,
3165                    vmoid: vmo_id.id,
3166                    length: 1,
3167                    ..Default::default()
3168                })
3169                .await
3170                .unwrap();
3171
3172            let mut response = BlockFifoResponse::default();
3173            reader.read_entries(&mut response).await.unwrap();
3174            assert_eq!(response.status, zx::sys::ZX_OK);
3175        });
3176
3177        // The response won't come back until we send on `tx`.
3178        assert!(fasync::TestExecutor::poll_until_stalled(&mut fut).await.is_pending());
3179
3180        let mut fut2 = pin!(proxy.get_volume_info());
3181
3182        // get_volume_info is set up to stall forever.
3183        assert!(fasync::TestExecutor::poll_until_stalled(&mut fut2).await.is_pending());
3184
3185        // If we now free up the first future, it should resolve; the stalled call to
3186        // get_volume_info should not block the fifo response.
3187        let _ = tx.send(());
3188
3189        assert!(fasync::TestExecutor::poll_until_stalled(&mut fut).await.is_ready());
3190    }
3191
3192    #[fuchsia::test]
3193    async fn test_request_flow_control() {
3194        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3195
3196        // The client will ensure that MAX_REQUESTS are queued up before firing `event`, and the
3197        // server will block until that happens.
3198        const MAX_REQUESTS: u64 = FIFO_MAX_REQUESTS as u64;
3199        let event = Arc::new((event_listener::Event::new(), AtomicBool::new(false)));
3200        let event_clone = event.clone();
3201        futures::join!(
3202            async move {
3203                let block_server = BlockServer::new(
3204                    BLOCK_SIZE,
3205                    Arc::new(MockInterface {
3206                        info: Some(DeviceInfo::Partition(PartitionInfo {
3207                            block_count: 1000,
3208                            ..Default::default()
3209                        })),
3210                        read_hook: Some(Box::new(move |_, _, _, _| {
3211                            let event_clone = event_clone.clone();
3212                            Box::pin(async move {
3213                                if !event_clone.1.load(Ordering::SeqCst) {
3214                                    event_clone.0.listen().await;
3215                                }
3216                                Ok(())
3217                            })
3218                        })),
3219                        ..MockInterface::default()
3220                    }),
3221                );
3222                block_server.handle_requests(stream).await.unwrap();
3223            },
3224            async move {
3225                let (session_proxy, server) = fidl::endpoints::create_proxy();
3226
3227                proxy.open_session(server).unwrap();
3228
3229                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
3230                let vmo_id = session_proxy
3231                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
3232                    .await
3233                    .unwrap()
3234                    .unwrap();
3235
3236                let mut fifo =
3237                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
3238                let (mut reader, mut writer) = fifo.async_io();
3239
3240                for i in 0..MAX_REQUESTS {
3241                    writer
3242                        .write_entries(&BlockFifoRequest {
3243                            command: BlockFifoCommand {
3244                                opcode: BlockOpcode::Read.into_primitive(),
3245                                ..Default::default()
3246                            },
3247                            reqid: (i + 1) as u32,
3248                            dev_offset: i,
3249                            vmoid: vmo_id.id,
3250                            length: 1,
3251                            ..Default::default()
3252                        })
3253                        .await
3254                        .unwrap();
3255                }
3256                assert!(
3257                    futures::poll!(pin!(writer.write_entries(&BlockFifoRequest {
3258                        command: BlockFifoCommand {
3259                            opcode: BlockOpcode::Read.into_primitive(),
3260                            ..Default::default()
3261                        },
3262                        reqid: u32::MAX,
3263                        dev_offset: MAX_REQUESTS,
3264                        vmoid: vmo_id.id,
3265                        length: 1,
3266                        ..Default::default()
3267                    })))
3268                    .is_pending()
3269                );
3270                // OK, let the server start to process.
3271                event.1.store(true, Ordering::SeqCst);
3272                event.0.notify(usize::MAX);
3273                // For each entry we read, make sure we can write a new one in.
3274                let mut finished_reqids = vec![];
3275                for i in MAX_REQUESTS..2 * MAX_REQUESTS {
3276                    let mut response = BlockFifoResponse::default();
3277                    reader.read_entries(&mut response).await.unwrap();
3278                    assert_eq!(response.status, zx::sys::ZX_OK);
3279                    finished_reqids.push(response.reqid);
3280                    writer
3281                        .write_entries(&BlockFifoRequest {
3282                            command: BlockFifoCommand {
3283                                opcode: BlockOpcode::Read.into_primitive(),
3284                                ..Default::default()
3285                            },
3286                            reqid: (i + 1) as u32,
3287                            dev_offset: i,
3288                            vmoid: vmo_id.id,
3289                            length: 1,
3290                            ..Default::default()
3291                        })
3292                        .await
3293                        .unwrap();
3294                }
3295                let mut response = BlockFifoResponse::default();
3296                for _ in 0..MAX_REQUESTS {
3297                    reader.read_entries(&mut response).await.unwrap();
3298                    assert_eq!(response.status, zx::sys::ZX_OK);
3299                    finished_reqids.push(response.reqid);
3300                }
3301                // Verify that we got a response for each request.  Note that we can't assume FIFO
3302                // ordering.
3303                finished_reqids.sort();
3304                assert_eq!(finished_reqids.len(), 2 * MAX_REQUESTS as usize);
3305                let mut i = 1;
3306                for reqid in finished_reqids {
3307                    assert_eq!(reqid, i);
3308                    i += 1;
3309                }
3310            }
3311        );
3312    }
3313
3314    #[fuchsia::test]
3315    async fn test_passthrough_io_with_fixed_map() {
3316        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3317
3318        let expected_op = Arc::new(Mutex::new(None));
3319        let expected_op_clone = expected_op.clone();
3320        futures::join!(
3321            async {
3322                let block_server = BlockServer::new(
3323                    BLOCK_SIZE,
3324                    Arc::new(IoMockInterface {
3325                        return_errors: false,
3326                        do_checks: true,
3327                        expected_op: expected_op_clone,
3328                    }),
3329                );
3330                block_server.handle_requests(stream).await.unwrap();
3331            },
3332            async move {
3333                let (session_proxy, server) = fidl::endpoints::create_proxy();
3334
3335                let mapping = fblock::BlockOffsetMapping { target_block_offset: 10, length: 20 };
3336                proxy.open_session_with_options(server, &[mapping]).unwrap();
3337
3338                let vmo = zx::Vmo::create(2 * zx::system_get_page_size() as u64).unwrap();
3339                let vmo_id = session_proxy
3340                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
3341                    .await
3342                    .unwrap()
3343                    .unwrap();
3344
3345                let mut fifo =
3346                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
3347                let (mut reader, mut writer) = fifo.async_io();
3348
3349                // READ
3350                *expected_op.lock() = Some(ExpectedOp::Read(11, 2, 3));
3351                writer
3352                    .write_entries(&BlockFifoRequest {
3353                        command: BlockFifoCommand {
3354                            opcode: BlockOpcode::Read.into_primitive(),
3355                            ..Default::default()
3356                        },
3357                        vmoid: vmo_id.id,
3358                        dev_offset: 1,
3359                        length: 2,
3360                        vmo_offset: 3,
3361                        ..Default::default()
3362                    })
3363                    .await
3364                    .unwrap();
3365
3366                let mut response = BlockFifoResponse::default();
3367                reader.read_entries(&mut response).await.unwrap();
3368                assert_eq!(response.status, zx::sys::ZX_OK);
3369
3370                // WRITE
3371                *expected_op.lock() = Some(ExpectedOp::Write(14, 5, 6));
3372                writer
3373                    .write_entries(&BlockFifoRequest {
3374                        command: BlockFifoCommand {
3375                            opcode: BlockOpcode::Write.into_primitive(),
3376                            ..Default::default()
3377                        },
3378                        vmoid: vmo_id.id,
3379                        dev_offset: 4,
3380                        length: 5,
3381                        vmo_offset: 6,
3382                        ..Default::default()
3383                    })
3384                    .await
3385                    .unwrap();
3386
3387                reader.read_entries(&mut response).await.unwrap();
3388                assert_eq!(response.status, zx::sys::ZX_OK);
3389
3390                // FLUSH
3391                *expected_op.lock() = Some(ExpectedOp::Flush);
3392                writer
3393                    .write_entries(&BlockFifoRequest {
3394                        command: BlockFifoCommand {
3395                            opcode: BlockOpcode::Flush.into_primitive(),
3396                            ..Default::default()
3397                        },
3398                        ..Default::default()
3399                    })
3400                    .await
3401                    .unwrap();
3402
3403                reader.read_entries(&mut response).await.unwrap();
3404                assert_eq!(response.status, zx::sys::ZX_OK);
3405
3406                // TRIM
3407                *expected_op.lock() = Some(ExpectedOp::Trim(17, 3));
3408                writer
3409                    .write_entries(&BlockFifoRequest {
3410                        command: BlockFifoCommand {
3411                            opcode: BlockOpcode::Trim.into_primitive(),
3412                            ..Default::default()
3413                        },
3414                        dev_offset: 7,
3415                        length: 3,
3416                        ..Default::default()
3417                    })
3418                    .await
3419                    .unwrap();
3420
3421                reader.read_entries(&mut response).await.unwrap();
3422                assert_eq!(response.status, zx::sys::ZX_OK);
3423
3424                // READ past window
3425                *expected_op.lock() = None;
3426                writer
3427                    .write_entries(&BlockFifoRequest {
3428                        command: BlockFifoCommand {
3429                            opcode: BlockOpcode::Read.into_primitive(),
3430                            ..Default::default()
3431                        },
3432                        vmoid: vmo_id.id,
3433                        dev_offset: 19,
3434                        length: 2,
3435                        vmo_offset: 3,
3436                        ..Default::default()
3437                    })
3438                    .await
3439                    .unwrap();
3440
3441                reader.read_entries(&mut response).await.unwrap();
3442                assert_eq!(response.status, zx::sys::ZX_ERR_OUT_OF_RANGE);
3443
3444                std::mem::drop(proxy);
3445            }
3446        );
3447    }
3448
3449    #[fuchsia::test]
3450    fn operation_map() {
3451        const BLOCK_SIZE: u32 = 512;
3452
3453        #[track_caller]
3454        fn expect_map_result(
3455            mut operation: Operation,
3456            mapping: Option<fblock::BlockOffsetMapping>,
3457            max_blocks: Option<NonZero<u32>>,
3458            expected_operations: Vec<Operation>,
3459        ) {
3460            let offset_map = mapping
3461                .map(|m| {
3462                    let map: BlockOffsetMapping = (&m).try_into().unwrap();
3463                    OffsetMap::new(vec![map]).unwrap()
3464                })
3465                .unwrap_or_else(OffsetMap::empty);
3466            let mut ops = vec![];
3467            while let Some(remainder) = operation.map(&offset_map, max_blocks, BLOCK_SIZE).unwrap()
3468            {
3469                ops.push(operation);
3470                operation = remainder;
3471            }
3472            ops.push(operation);
3473            assert_eq!(ops, expected_operations);
3474        }
3475
3476        // No limits
3477        expect_map_result(
3478            Operation::Read {
3479                device_block_offset: 10,
3480                block_count: 200,
3481                _unused: 0,
3482                vmo_offset: 0,
3483                options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3484            },
3485            None,
3486            None,
3487            vec![Operation::Read {
3488                device_block_offset: 10,
3489                block_count: 200,
3490                _unused: 0,
3491                vmo_offset: 0,
3492                options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3493            }],
3494        );
3495
3496        // Max block count
3497        expect_map_result(
3498            Operation::Read {
3499                device_block_offset: 10,
3500                block_count: 200,
3501                _unused: 0,
3502                vmo_offset: 0,
3503                options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3504            },
3505            None,
3506            NonZero::new(120),
3507            vec![
3508                Operation::Read {
3509                    device_block_offset: 10,
3510                    block_count: 120,
3511                    _unused: 0,
3512                    vmo_offset: 0,
3513                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3514                },
3515                Operation::Read {
3516                    device_block_offset: 130,
3517                    block_count: 80,
3518                    _unused: 0,
3519                    vmo_offset: 120 * BLOCK_SIZE as u64,
3520                    options: ReadOptions {
3521                        // The DUN should be offset by the number of blocks in the first request.
3522                        inline_crypto: InlineCryptoOptions::enabled(1, 1000 + 120),
3523                    },
3524                },
3525            ],
3526        );
3527        expect_map_result(
3528            Operation::Trim { device_block_offset: 10, block_count: 200 },
3529            None,
3530            NonZero::new(120),
3531            vec![Operation::Trim { device_block_offset: 10, block_count: 200 }],
3532        );
3533
3534        // Remapping + Max block count
3535        expect_map_result(
3536            Operation::Read {
3537                device_block_offset: 0,
3538                block_count: 200,
3539                _unused: 0,
3540                vmo_offset: 0,
3541                options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3542            },
3543            Some(fblock::BlockOffsetMapping { target_block_offset: 100, length: 200 }),
3544            NonZero::new(120),
3545            vec![
3546                Operation::Read {
3547                    device_block_offset: 100,
3548                    block_count: 120,
3549                    _unused: 0,
3550                    vmo_offset: 0,
3551                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3552                },
3553                Operation::Read {
3554                    device_block_offset: 220,
3555                    block_count: 80,
3556                    _unused: 0,
3557                    vmo_offset: 120 * BLOCK_SIZE as u64,
3558                    options: ReadOptions {
3559                        inline_crypto: InlineCryptoOptions::enabled(1, 1000 + 120),
3560                    },
3561                },
3562            ],
3563        );
3564        expect_map_result(
3565            Operation::Trim { device_block_offset: 0, block_count: 200 },
3566            Some(fblock::BlockOffsetMapping { target_block_offset: 100, length: 200 }),
3567            NonZero::new(120),
3568            vec![Operation::Trim { device_block_offset: 100, block_count: 200 }],
3569        );
3570
3571        // Multi-extent remapping
3572        let multi_extent_map = OffsetMap::new(vec![
3573            BlockOffsetMapping { target_block_offset: 100, length: 10 },
3574            BlockOffsetMapping { target_block_offset: 200, length: 20 },
3575        ])
3576        .unwrap();
3577
3578        fn expect_multi_map_result(
3579            mut operation: Operation,
3580            offset_map: &OffsetMap,
3581            max_blocks: Option<NonZero<u32>>,
3582            expected_operations: Vec<Operation>,
3583        ) {
3584            let mut ops = vec![];
3585            while let Some(remainder) = operation.map(offset_map, max_blocks, BLOCK_SIZE).unwrap() {
3586                ops.push(operation);
3587                operation = remainder;
3588            }
3589            ops.push(operation);
3590            assert_eq!(ops, expected_operations);
3591        }
3592
3593        // Read spanning multi-extents (logical offset 5, count 15 -> 5 blocks in extent 0, 10 in
3594        // extent 1)
3595        expect_multi_map_result(
3596            Operation::Read {
3597                device_block_offset: 5,
3598                block_count: 15,
3599                _unused: 0,
3600                vmo_offset: 0,
3601                options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3602            },
3603            &multi_extent_map,
3604            None,
3605            vec![
3606                Operation::Read {
3607                    device_block_offset: 105,
3608                    block_count: 5,
3609                    _unused: 0,
3610                    vmo_offset: 0,
3611                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3612                },
3613                Operation::Read {
3614                    device_block_offset: 200,
3615                    block_count: 10,
3616                    _unused: 0,
3617                    vmo_offset: 5 * BLOCK_SIZE as u64,
3618                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1005) },
3619                },
3620            ],
3621        );
3622
3623        // Read spanning multi-extents with max_transfer_blocks limit
3624        expect_multi_map_result(
3625            Operation::Read {
3626                device_block_offset: 5,
3627                block_count: 15,
3628                _unused: 0,
3629                vmo_offset: 0,
3630                options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3631            },
3632            &multi_extent_map,
3633            NonZero::new(7),
3634            vec![
3635                Operation::Read {
3636                    device_block_offset: 105,
3637                    block_count: 5,
3638                    _unused: 0,
3639                    vmo_offset: 0,
3640                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1000) },
3641                },
3642                Operation::Read {
3643                    device_block_offset: 200,
3644                    block_count: 7,
3645                    _unused: 0,
3646                    vmo_offset: 5 * BLOCK_SIZE as u64,
3647                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1005) },
3648                },
3649                Operation::Read {
3650                    device_block_offset: 207,
3651                    block_count: 3,
3652                    _unused: 0,
3653                    vmo_offset: 12 * BLOCK_SIZE as u64,
3654                    options: ReadOptions { inline_crypto: InlineCryptoOptions::enabled(1, 1012) },
3655                },
3656            ],
3657        );
3658
3659        // Write spanning multi-extents
3660        expect_multi_map_result(
3661            Operation::Write {
3662                device_block_offset: 5,
3663                block_count: 15,
3664                _unused: 0,
3665                vmo_offset: 0,
3666                options: WriteOptions {
3667                    inline_crypto: InlineCryptoOptions::enabled(1, 2000),
3668                    flags: WriteFlags::empty(),
3669                },
3670            },
3671            &multi_extent_map,
3672            None,
3673            vec![
3674                Operation::Write {
3675                    device_block_offset: 105,
3676                    block_count: 5,
3677                    _unused: 0,
3678                    vmo_offset: 0,
3679                    options: WriteOptions {
3680                        inline_crypto: InlineCryptoOptions::enabled(1, 2000),
3681                        flags: WriteFlags::empty(),
3682                    },
3683                },
3684                Operation::Write {
3685                    device_block_offset: 200,
3686                    block_count: 10,
3687                    _unused: 0,
3688                    vmo_offset: 5 * BLOCK_SIZE as u64,
3689                    options: WriteOptions {
3690                        inline_crypto: InlineCryptoOptions::enabled(1, 2005),
3691                        flags: WriteFlags::empty(),
3692                    },
3693                },
3694            ],
3695        );
3696
3697        // Trim spanning multi-extents
3698        expect_multi_map_result(
3699            Operation::Trim { device_block_offset: 5, block_count: 15 },
3700            &multi_extent_map,
3701            None,
3702            vec![
3703                Operation::Trim { device_block_offset: 105, block_count: 5 },
3704                Operation::Trim { device_block_offset: 200, block_count: 10 },
3705            ],
3706        );
3707
3708        // Large extent test (length > u32::MAX)
3709        let large_extent_map = OffsetMap::new(vec![BlockOffsetMapping {
3710            target_block_offset: 100,
3711            length: (u32::MAX as u64) + 10,
3712        }])
3713        .unwrap();
3714
3715        assert_eq!(large_extent_map.map(0), Some((100, u32::MAX)));
3716    }
3717
3718    // Verifies that if the pre-flush (for a simulated barrier) fails, the write is not executed.
3719    #[fuchsia::test]
3720    async fn test_pre_barrier_flush_failure() {
3721        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3722
3723        struct NoBarrierInterface;
3724        impl super::async_interface::Interface for NoBarrierInterface {
3725            fn get_info(&self) -> Cow<'_, DeviceInfo> {
3726                Cow::Owned(DeviceInfo::Partition(PartitionInfo {
3727                    device_flags: fblock::DeviceFlag::empty(), // No BARRIER_SUPPORT
3728                    max_transfer_blocks: NonZero::new(100),
3729                    start_block_offset: Some(0),
3730                    block_count: 100,
3731                    type_guid: [0; 16],
3732                    instance_guid: [0; 16],
3733                    name: "test".to_string(),
3734                    flags: Some(0),
3735                }))
3736            }
3737            async fn on_attach_vmo(&self, _vmo: &zx::Vmo) -> Result<(), zx::Status> {
3738                Ok(())
3739            }
3740            async fn read(
3741                &self,
3742                _: u64,
3743                _: u32,
3744                _: &Arc<zx::Vmo>,
3745                _: u64,
3746                _: ReadOptions,
3747                _: TraceFlowId,
3748            ) -> Result<(), zx::Status> {
3749                unreachable!()
3750            }
3751            async fn write(
3752                &self,
3753                _: u64,
3754                _: u32,
3755                _: &Arc<zx::Vmo>,
3756                _: u64,
3757                _: WriteOptions,
3758                _: TraceFlowId,
3759            ) -> Result<(), zx::Status> {
3760                panic!("Write should not be called");
3761            }
3762            async fn flush(&self, _: TraceFlowId) -> Result<(), zx::Status> {
3763                Err(zx::Status::IO)
3764            }
3765            async fn trim(&self, _: u64, _: u32, _: TraceFlowId) -> Result<(), zx::Status> {
3766                unreachable!()
3767            }
3768        }
3769
3770        futures::join!(
3771            async move {
3772                let block_server = BlockServer::new(BLOCK_SIZE, Arc::new(NoBarrierInterface));
3773                block_server.handle_requests(stream).await.unwrap();
3774            },
3775            async move {
3776                let (session_proxy, server) = fidl::endpoints::create_proxy();
3777                proxy.open_session(server).unwrap();
3778                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
3779                let vmo_id = session_proxy
3780                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
3781                    .await
3782                    .unwrap()
3783                    .unwrap();
3784
3785                let mut fifo =
3786                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
3787                let (mut reader, mut writer) = fifo.async_io();
3788
3789                writer
3790                    .write_entries(&BlockFifoRequest {
3791                        command: BlockFifoCommand {
3792                            opcode: BlockOpcode::Write.into_primitive(),
3793                            flags: BlockIoFlag::PRE_BARRIER.bits(),
3794                            ..Default::default()
3795                        },
3796                        vmoid: vmo_id.id,
3797                        length: 1,
3798                        ..Default::default()
3799                    })
3800                    .await
3801                    .unwrap();
3802
3803                let mut response = BlockFifoResponse::default();
3804                reader.read_entries(&mut response).await.unwrap();
3805                assert_eq!(response.status, zx::sys::ZX_ERR_IO);
3806            }
3807        );
3808    }
3809
3810    // Verifies that if the write fails when a post-flush is required (for a simulated FUA), the
3811    // post-flush is not executed.
3812    #[fuchsia::test]
3813    async fn test_post_barrier_write_failure() {
3814        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3815
3816        struct NoBarrierInterface;
3817        impl super::async_interface::Interface for NoBarrierInterface {
3818            fn get_info(&self) -> Cow<'_, DeviceInfo> {
3819                Cow::Owned(DeviceInfo::Partition(PartitionInfo {
3820                    device_flags: fblock::DeviceFlag::empty(), // No FUA_SUPPORT
3821                    max_transfer_blocks: NonZero::new(100),
3822                    start_block_offset: Some(0),
3823                    block_count: 100,
3824                    type_guid: [0; 16],
3825                    instance_guid: [0; 16],
3826                    name: "test".to_string(),
3827                    flags: Some(0),
3828                }))
3829            }
3830            async fn on_attach_vmo(&self, _vmo: &zx::Vmo) -> Result<(), zx::Status> {
3831                Ok(())
3832            }
3833            async fn read(
3834                &self,
3835                _: u64,
3836                _: u32,
3837                _: &Arc<zx::Vmo>,
3838                _: u64,
3839                _: ReadOptions,
3840                _: TraceFlowId,
3841            ) -> Result<(), zx::Status> {
3842                unreachable!()
3843            }
3844            async fn write(
3845                &self,
3846                _: u64,
3847                _: u32,
3848                _: &Arc<zx::Vmo>,
3849                _: u64,
3850                _: WriteOptions,
3851                _: TraceFlowId,
3852            ) -> Result<(), zx::Status> {
3853                Err(zx::Status::IO)
3854            }
3855            async fn flush(&self, _: TraceFlowId) -> Result<(), zx::Status> {
3856                panic!("Flush should not be called")
3857            }
3858            async fn trim(&self, _: u64, _: u32, _: TraceFlowId) -> Result<(), zx::Status> {
3859                unreachable!()
3860            }
3861        }
3862
3863        futures::join!(
3864            async move {
3865                let block_server = BlockServer::new(BLOCK_SIZE, Arc::new(NoBarrierInterface));
3866                block_server.handle_requests(stream).await.unwrap();
3867            },
3868            async move {
3869                let (session_proxy, server) = fidl::endpoints::create_proxy();
3870                proxy.open_session(server).unwrap();
3871                let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
3872                let vmo_id = session_proxy
3873                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
3874                    .await
3875                    .unwrap()
3876                    .unwrap();
3877
3878                let mut fifo =
3879                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
3880                let (mut reader, mut writer) = fifo.async_io();
3881
3882                writer
3883                    .write_entries(&BlockFifoRequest {
3884                        command: BlockFifoCommand {
3885                            opcode: BlockOpcode::Write.into_primitive(),
3886                            flags: BlockIoFlag::FORCE_ACCESS.bits(),
3887                            ..Default::default()
3888                        },
3889                        vmoid: vmo_id.id,
3890                        length: 1,
3891                        ..Default::default()
3892                    })
3893                    .await
3894                    .unwrap();
3895
3896                let mut response = BlockFifoResponse::default();
3897                reader.read_entries(&mut response).await.unwrap();
3898                assert_eq!(response.status, zx::sys::ZX_ERR_IO);
3899            }
3900        );
3901    }
3902
3903    /// Verifies that group IDs are isolated per session.
3904    ///
3905    /// Even if two independent sessions on the same BlockServer use the same group ID,
3906    /// their in-flight transaction groups must remain isolated.
3907    #[fuchsia::test]
3908    async fn test_group_ids_isolated_per_session() {
3909        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
3910
3911        futures::join!(
3912            async {
3913                let block_server = BlockServer::new(
3914                    BLOCK_SIZE,
3915                    // MockInterface::flush() is a no-op that returns Ok(()).
3916                    Arc::new(MockInterface::default()),
3917                );
3918                block_server.handle_requests(stream).await.unwrap();
3919            },
3920            async move {
3921                async fn settle() {
3922                    // Let the single-threaded executor drain server work.
3923                    for _ in 0..32 {
3924                        fasync::yield_now().await;
3925                    }
3926                }
3927
3928                // --- Open session A. ---
3929                let (session_a, server_a) = fidl::endpoints::create_proxy();
3930                proxy.open_session(server_a).unwrap();
3931                let mut fifo_a =
3932                    fasync::Fifo::from_fifo(session_a.get_fifo().await.unwrap().unwrap());
3933
3934                // --- Open session B. ---
3935                let (session_b, server_b) = fidl::endpoints::create_proxy();
3936                proxy.open_session(server_b).unwrap();
3937                let mut fifo_b =
3938                    fasync::Fifo::from_fifo(session_b.get_fifo().await.unwrap().unwrap());
3939
3940                // ----------------------------------------------------------------
3941                // Control: with no interference, A's two-part Flush group is OK.
3942                // ----------------------------------------------------------------
3943                {
3944                    let (mut reader_a, mut writer_a) = fifo_a.async_io();
3945                    writer_a
3946                        .write_entries(&BlockFifoRequest {
3947                            command: BlockFifoCommand {
3948                                opcode: BlockOpcode::Flush.into_primitive(),
3949                                flags: BlockIoFlag::GROUP_ITEM.bits(),
3950                                ..Default::default()
3951                            },
3952                            group: 1,
3953                            reqid: 0xAAAA,
3954                            ..Default::default()
3955                        })
3956                        .await
3957                        .unwrap();
3958                    settle().await;
3959                    writer_a
3960                        .write_entries(&BlockFifoRequest {
3961                            command: BlockFifoCommand {
3962                                opcode: BlockOpcode::Flush.into_primitive(),
3963                                flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
3964                                ..Default::default()
3965                            },
3966                            group: 1,
3967                            reqid: 0xAAAA,
3968                            ..Default::default()
3969                        })
3970                        .await
3971                        .unwrap();
3972                    let mut response = BlockFifoResponse::default();
3973                    reader_a.read_entries(&mut response).await.unwrap();
3974                    assert_eq!(response.reqid, 0xAAAA);
3975                    assert_eq!(
3976                        response.status,
3977                        zx::sys::ZX_OK,
3978                        "control: A's valid Flush group must succeed"
3979                    );
3980                }
3981
3982                // ----------------------------------------------------------------
3983                // Run concurrent group requests with the same group ID (7) on
3984                // both sessions, and verify they both succeed independently.
3985                // ----------------------------------------------------------------
3986
3987                // Step 1: Session A starts group 7.
3988                {
3989                    let (_reader_a, mut writer_a) = fifo_a.async_io();
3990                    writer_a
3991                        .write_entries(&BlockFifoRequest {
3992                            command: BlockFifoCommand {
3993                                opcode: BlockOpcode::Flush.into_primitive(),
3994                                flags: BlockIoFlag::GROUP_ITEM.bits(),
3995                                ..Default::default()
3996                            },
3997                            group: 7,
3998                            reqid: 100,
3999                            ..Default::default()
4000                        })
4001                        .await
4002                        .unwrap();
4003                }
4004                settle().await;
4005
4006                // Step 2: Session B starts group 7.
4007                {
4008                    let (_reader_b, mut writer_b) = fifo_b.async_io();
4009                    writer_b
4010                        .write_entries(&BlockFifoRequest {
4011                            command: BlockFifoCommand {
4012                                opcode: BlockOpcode::Flush.into_primitive(),
4013                                flags: BlockIoFlag::GROUP_ITEM.bits(),
4014                                ..Default::default()
4015                            },
4016                            group: 7,
4017                            reqid: 200,
4018                            ..Default::default()
4019                        })
4020                        .await
4021                        .unwrap();
4022                }
4023                settle().await;
4024
4025                // Step 3: Session A finishes group 7.
4026                {
4027                    let (_reader_a, mut writer_a) = fifo_a.async_io();
4028                    writer_a
4029                        .write_entries(&BlockFifoRequest {
4030                            command: BlockFifoCommand {
4031                                opcode: BlockOpcode::Flush.into_primitive(),
4032                                flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
4033                                ..Default::default()
4034                            },
4035                            group: 7,
4036                            reqid: 100,
4037                            ..Default::default()
4038                        })
4039                        .await
4040                        .unwrap();
4041                }
4042                settle().await;
4043
4044                // Step 4: Session B finishes group 7.
4045                {
4046                    let (_reader_b, mut writer_b) = fifo_b.async_io();
4047                    writer_b
4048                        .write_entries(&BlockFifoRequest {
4049                            command: BlockFifoCommand {
4050                                opcode: BlockOpcode::Flush.into_primitive(),
4051                                flags: (BlockIoFlag::GROUP_ITEM | BlockIoFlag::GROUP_LAST).bits(),
4052                                ..Default::default()
4053                            },
4054                            group: 7,
4055                            reqid: 200,
4056                            ..Default::default()
4057                        })
4058                        .await
4059                        .unwrap();
4060                }
4061                settle().await;
4062
4063                // Verify Session A's response.
4064                {
4065                    let (mut reader_a, _writer_a) = fifo_a.async_io();
4066                    let mut response_a = BlockFifoResponse::default();
4067                    reader_a.read_entries(&mut response_a).await.unwrap();
4068                    assert_eq!(response_a.reqid, 100);
4069                    assert_eq!(response_a.group, 7);
4070                    assert_eq!(response_a.status, zx::sys::ZX_OK);
4071                }
4072
4073                // Verify Session B's response.
4074                {
4075                    let (mut reader_b, _writer_b) = fifo_b.async_io();
4076                    let mut response_b = BlockFifoResponse::default();
4077                    reader_b.read_entries(&mut response_b).await.unwrap();
4078                    assert_eq!(response_b.reqid, 200);
4079                    assert_eq!(response_b.group, 7);
4080                    assert_eq!(response_b.status, zx::sys::ZX_OK);
4081                }
4082
4083                std::mem::drop(session_a);
4084                std::mem::drop(session_b);
4085                std::mem::drop(proxy);
4086            }
4087        );
4088    }
4089
4090    #[fuchsia::test]
4091    async fn test_unmapped_request_out_of_range() {
4092        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
4093        let _server = fasync::Task::spawn(async move {
4094            let block_server = BlockServer::new(4096, Arc::new(MockInterface::default()));
4095            let _ = block_server.handle_requests(stream).await;
4096        });
4097
4098        let (session_proxy, server) = fidl::endpoints::create_proxy();
4099        proxy.open_session(server).unwrap();
4100
4101        let vmo = zx::Vmo::create(zx::system_get_page_size() as u64).unwrap();
4102        let vmo_id = session_proxy
4103            .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
4104            .await
4105            .unwrap()
4106            .unwrap();
4107
4108        let mut fifo = fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
4109        let (mut reader, mut writer) = fifo.async_io();
4110
4111        // Attempting to read at offset 100 on a 100-block unmapped partition should fail with
4112        // OUT_OF_RANGE.
4113        writer
4114            .write_entries(&BlockFifoRequest {
4115                command: BlockFifoCommand {
4116                    opcode: BlockOpcode::Read.into_primitive(),
4117                    ..Default::default()
4118                },
4119                reqid: 1,
4120                vmoid: vmo_id.id,
4121                length: 1,
4122                dev_offset: 100,
4123                ..Default::default()
4124            })
4125            .await
4126            .unwrap();
4127
4128        let mut response = BlockFifoResponse::default();
4129        reader.read_entries(&mut response).await.unwrap();
4130        assert_eq!(zx::Status::ok(response.status), Err(zx::Status::OUT_OF_RANGE));
4131    }
4132
4133    #[test]
4134    fn test_offset_map_coalescing() {
4135        use crate::{BlockOffsetMapping, OffsetMap};
4136
4137        // Mappings 0 and 1 are contiguous in device offset (100..150 + 150..180 = 100..180).
4138        // Mapping 2 is not contiguous with 1 (target 250 instead of 180).
4139        let mappings = vec![
4140            BlockOffsetMapping { target_block_offset: 100, length: 50 },
4141            BlockOffsetMapping { target_block_offset: 150, length: 30 },
4142            BlockOffsetMapping { target_block_offset: 250, length: 20 },
4143        ];
4144        let map = OffsetMap::new(crate::coalesce_mappings(mappings)).unwrap();
4145
4146        assert_eq!(map.mappings().len(), 2);
4147        assert_eq!(map.mappings()[0].target_block_offset, 100);
4148        assert_eq!(map.mappings()[0].length, 80);
4149        assert_eq!(map.mappings()[1].target_block_offset, 250);
4150        assert_eq!(map.mappings()[1].length, 20);
4151
4152        // Logical offset 10 falls in the first extent, but since it coalesces with the second,
4153        // extent_remaining_blocks extends to logical offset 80 (80 - 10 = 70).
4154        assert_eq!(map.map(10), Some((110, 70)));
4155        // Logical offset 55 falls in the second extent, remaining blocks up to logical 80
4156        // (80 - 55 = 25).
4157        assert_eq!(map.map(55), Some((155, 25)));
4158        // Logical offset 85 falls in the third extent, which is not coalesced with the second
4159        // (100 - 85 = 15).
4160        assert_eq!(map.map(85), Some((255, 15)));
4161    }
4162    #[fuchsia::test]
4163    async fn test_open_session_with_options_errors() {
4164        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
4165
4166        futures::join!(
4167            async {
4168                let block_server = BlockServer::new(BLOCK_SIZE, Arc::new(MockInterface::default()));
4169                let _ = block_server.handle_requests(stream).await;
4170            },
4171            async move {
4172                // Test 1: Mapping with zero length -> should close session with
4173                // INVALID_ARGS epitaph.
4174                {
4175                    let (session_proxy, server) = fidl::endpoints::create_proxy();
4176                    proxy
4177                        .open_session_with_options(
4178                            server,
4179                            &[fblock::BlockOffsetMapping { target_block_offset: 0, length: 0 }],
4180                        )
4181                        .unwrap();
4182                    let res = session_proxy.get_fifo().await;
4183                    assert_matches!(
4184                        res,
4185                        Err(fidl::Error::ClientChannelClosed { epitaph, .. })
4186                            if epitaph == zx::Status::INVALID_ARGS
4187                    );
4188                }
4189
4190                // Test 2: Mappings out of range (exceeding block_count) -> should close with
4191                // OUT_OF_RANGE epitaph.
4192                {
4193                    let (session_proxy, server) = fidl::endpoints::create_proxy();
4194                    let mapping = fblock::BlockOffsetMapping {
4195                        target_block_offset: u64::MAX - 10, // Way past block_count
4196                        length: 10,
4197                    };
4198                    proxy.open_session_with_options(server, &[mapping]).unwrap();
4199                    let res = session_proxy.get_fifo().await;
4200                    assert_matches!(
4201                        res,
4202                        Err(fidl::Error::ClientChannelClosed { epitaph, .. })
4203                            if epitaph == zx::Status::OUT_OF_RANGE
4204                    );
4205                }
4206            }
4207        );
4208    }
4209
4210    #[fuchsia::test]
4211    async fn test_split_request_failure_aborts_subsequent_chunks() {
4212        let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<fblock::BlockMarker>();
4213
4214        let read_calls = Arc::new(AtomicU64::new(0));
4215        let read_calls_clone = read_calls.clone();
4216
4217        struct MockSplitInterface {
4218            read_calls: Arc<AtomicU64>,
4219        }
4220
4221        impl super::async_interface::Interface for MockSplitInterface {
4222            fn get_info(&self) -> Cow<'_, DeviceInfo> {
4223                Cow::Owned(DeviceInfo::Partition(PartitionInfo {
4224                    device_flags: fblock::DeviceFlag::READONLY,
4225                    max_transfer_blocks: NonZero::new(5),
4226                    start_block_offset: Some(0),
4227                    block_count: 100,
4228                    type_guid: [1; 16],
4229                    instance_guid: [2; 16],
4230                    name: "foo".to_string(),
4231                    flags: Some(0),
4232                }))
4233            }
4234
4235            async fn on_attach_vmo(&self, _vmo: &zx::Vmo) -> Result<(), zx::Status> {
4236                Ok(())
4237            }
4238
4239            async fn read(
4240                &self,
4241                _device_block_offset: u64,
4242                _block_count: u32,
4243                _vmo: &Arc<zx::Vmo>,
4244                _vmo_offset: u64,
4245                _opts: ReadOptions,
4246                _trace_flow_id: TraceFlowId,
4247            ) -> Result<(), zx::Status> {
4248                let call_num = self.read_calls.fetch_add(1, Ordering::Relaxed);
4249                if call_num == 0 { Err(zx::Status::IO) } else { Ok(()) }
4250            }
4251
4252            async fn write(
4253                &self,
4254                _device_block_offset: u64,
4255                _block_count: u32,
4256                _vmo: &Arc<zx::Vmo>,
4257                _vmo_offset: u64,
4258                _opts: WriteOptions,
4259                _trace_flow_id: TraceFlowId,
4260            ) -> Result<(), zx::Status> {
4261                unreachable!()
4262            }
4263
4264            async fn flush(&self, _trace_flow_id: TraceFlowId) -> Result<(), zx::Status> {
4265                Ok(())
4266            }
4267
4268            async fn trim(
4269                &self,
4270                _device_block_offset: u64,
4271                _block_count: u32,
4272                _trace_flow_id: TraceFlowId,
4273            ) -> Result<(), zx::Status> {
4274                unreachable!()
4275            }
4276        }
4277
4278        futures::join!(
4279            async move {
4280                let block_server = BlockServer::new(
4281                    BLOCK_SIZE,
4282                    Arc::new(MockSplitInterface { read_calls: read_calls_clone }),
4283                );
4284                let _ = block_server.handle_requests(stream).await;
4285            },
4286            async move {
4287                let (session_proxy, server) = fidl::endpoints::create_proxy();
4288                proxy.open_session(server).unwrap();
4289
4290                let vmo = zx::Vmo::create(2 * zx::system_get_page_size() as u64).unwrap();
4291                let vmo_id = session_proxy
4292                    .attach_vmo(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap())
4293                    .await
4294                    .unwrap()
4295                    .unwrap();
4296
4297                let mut fifo =
4298                    fasync::Fifo::from_fifo(session_proxy.get_fifo().await.unwrap().unwrap());
4299                let (mut reader, mut writer) = fifo.async_io();
4300
4301                // Send a request for 10 blocks. With max_transfer_blocks = 5, this will be split
4302                // into 2 chunks of 5 blocks each.
4303                writer
4304                    .write_entries(&BlockFifoRequest {
4305                        command: BlockFifoCommand {
4306                            opcode: BlockOpcode::Read.into_primitive(),
4307                            ..Default::default()
4308                        },
4309                        vmoid: vmo_id.id,
4310                        length: 10,
4311                        dev_offset: 0,
4312                        reqid: 1,
4313                        ..Default::default()
4314                    })
4315                    .await
4316                    .unwrap();
4317
4318                let mut response = BlockFifoResponse::default();
4319                reader.read_entries(&mut response).await.unwrap();
4320                assert_ne!(response.status, zx::sys::ZX_OK);
4321
4322                // Verify that only the first chunk was submitted to `read`. The second chunk
4323                // should have been aborted when `map_request` saw active_request.status != OK.
4324                assert_eq!(read_calls.load(Ordering::Relaxed), 1);
4325
4326                std::mem::drop(proxy);
4327            }
4328        );
4329    }
4330
4331    #[fuchsia::test]
4332    async fn test_mapper_open_session() {
4333        use crate::callback_interface::SessionManager;
4334        use crate::testing::MockInterface;
4335
4336        let (tx, _rx) = std::sync::mpsc::channel();
4337        let interface = Arc::new(MockInterface::new(tx));
4338        let session_manager = Arc::new(SessionManager::new(interface, 512));
4339        let block_server = BlockServer::new(512, session_manager);
4340
4341        let (mapper_proxy, mapper_stream) =
4342            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
4343        let scope = fasync::Scope::new();
4344        scope.spawn(async move {
4345            let _ = block_server.handle_mapper_requests(mapper_stream).await;
4346        });
4347
4348        // 1. Both port and delivery queue provided (with pager):
4349        let (_mapper_session_proxy, mapper_session_server) =
4350            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4351        let mapping_vmo = zx::Vmo::create(4096).unwrap();
4352        let port = zx::Port::create();
4353        let delivery_queue = zx::Vmo::create(4096).unwrap();
4354        let res = mapper_proxy
4355            .open_session(mapper_session_server, mapping_vmo, Some(port), Some(delivery_queue))
4356            .await
4357            .unwrap();
4358        assert_matches!(res, Ok(()));
4359
4360        // 2. Neither port nor delivery queue provided (pager-less):
4361        let (_mapper_session_proxy, mapper_session_server) =
4362            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4363        let mapping_vmo = zx::Vmo::create(4096).unwrap();
4364        let res = mapper_proxy
4365            .open_session(mapper_session_server, mapping_vmo, None, None)
4366            .await
4367            .unwrap();
4368        assert_matches!(res, Ok(()));
4369
4370        // 3. Port provided without delivery queue (invalid):
4371        let (_mapper_session_proxy, mapper_session_server) =
4372            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4373        let mapping_vmo = zx::Vmo::create(4096).unwrap();
4374        let port = zx::Port::create();
4375        let res = mapper_proxy
4376            .open_session(mapper_session_server, mapping_vmo, Some(port), None)
4377            .await
4378            .unwrap();
4379        assert_eq!(res, Err(zx::sys::ZX_ERR_INVALID_ARGS));
4380
4381        // 4. Delivery queue provided without port (invalid):
4382        let (_mapper_session_proxy, mapper_session_server) =
4383            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4384        let mapping_vmo = zx::Vmo::create(4096).unwrap();
4385        let delivery_queue = zx::Vmo::create(4096).unwrap();
4386        let res = mapper_proxy
4387            .open_session(mapper_session_server, mapping_vmo, None, Some(delivery_queue))
4388            .await
4389            .unwrap();
4390        assert_eq!(res, Err(zx::sys::ZX_ERR_INVALID_ARGS));
4391    }
4392
4393    #[fuchsia::test]
4394    async fn test_mapper_blob_page_request() {
4395        use crate::callback_interface::SessionManager;
4396        use crate::testing::MockInterface;
4397
4398        let (tx, rx) = std::sync::mpsc::channel();
4399        let interface = Arc::new(MockInterface::new(tx));
4400        let session_manager = Arc::new(SessionManager::new(interface.clone(), 512));
4401        let block_server = BlockServer::new(512, session_manager.clone());
4402
4403        let sm_completer = session_manager.clone();
4404        std::thread::spawn(move || {
4405            while let Ok(req) = rx.recv() {
4406                if let Operation::Read { vmo_offset, block_count, .. } = req.operation {
4407                    if let Some(vmo) = req.vmo {
4408                        let data = vec![0xABu8; (block_count * 512) as usize];
4409                        vmo.write(&data, vmo_offset).unwrap();
4410                    }
4411                }
4412                sm_completer.complete_request(req.request_id, Ok(()));
4413            }
4414        });
4415
4416        let (mapper_proxy, mapper_stream) =
4417            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
4418        let scope = fasync::Scope::new();
4419        scope.spawn(async move {
4420            let _ = block_server.handle_mapper_requests(mapper_stream).await;
4421        });
4422
4423        let pager = Arc::new(zx::Pager::create(zx::PagerOptions::empty()).unwrap());
4424        let port = zx::Port::create();
4425        let key = 1001u64;
4426        let paged_vmo = pager.create_vmo(zx::VmoOptions::empty(), &port, key, 4096).unwrap();
4427
4428        let (_mapper_session_proxy, mapper_session_server) =
4429            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4430        let mapping_vmo = zx::Vmo::create(65536).unwrap();
4431        let delivery_queue = zx::Vmo::create(65536).unwrap();
4432        let vmo_provider = Arc::new(blob_pager_and_verifier::TestVmoProvider::new(
4433            pager.clone(),
4434            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4435        ));
4436        vmo_provider
4437            .register_vmo(key, paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap());
4438        let receiver = vmo_fifo::Receiver::<mapping::RawDeliveryCommand>::new(
4439            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4440            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
4441        )
4442        .unwrap();
4443        let _delivery_processor = blob_pager_and_verifier::DeliveryQueueProcessor::spawn(
4444            receiver,
4445            vmo_provider,
4446            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4447        )
4448        .unwrap();
4449
4450        let extents =
4451            mapping::Extents::try_new([mapping::Extent::new(0..4096, Some(0))], 0).unwrap();
4452        let data_extent_words: Vec<u64> = mapping::Extents::encode_extents(&extents).collect();
4453        let mut payload_bytes = Vec::new();
4454        for w in &data_extent_words {
4455            payload_bytes.extend_from_slice(&w.to_le_bytes());
4456        }
4457
4458        let cmd = mapping::RawMappingCommand {
4459            opcode: mapping::MAPPINGS_COMMAND,
4460            offset: 0,
4461            key,
4462            stored_size: 4096,
4463            device_offset: 0,
4464            metadata_count: 0,
4465            extent_count: data_extent_words.len() as u32,
4466        };
4467
4468        let mut sender = vmo_fifo::SyncSender::<mapping::RawMappingCommand>::new(
4469            mapping_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4470            1024,
4471            256,
4472        )
4473        .unwrap();
4474        let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
4475        payload_buf.data().copy_from_slice(&payload_bytes);
4476        payload_buf.commit(cmd).unwrap();
4477
4478        let res = mapper_proxy
4479            .open_session(mapper_session_server, mapping_vmo, Some(port), Some(delivery_queue))
4480            .await
4481            .unwrap();
4482        assert_matches!(res, Ok(()));
4483
4484        let reader_thread = std::thread::spawn(move || {
4485            let mut buf = [0u8; 4096];
4486            paged_vmo.read(&mut buf, 0).expect("paged vmo read failed");
4487            buf
4488        });
4489
4490        let read_bytes = reader_thread.join().unwrap();
4491        assert_eq!(read_bytes.len(), 4096);
4492        assert_eq!(read_bytes, [0xABu8; 4096]);
4493    }
4494
4495    #[fuchsia::test]
4496    async fn test_mapper_hierarchical_session() {
4497        use crate::callback_interface::SessionManager;
4498        use crate::testing::MockInterface;
4499
4500        let (tx, rx) = std::sync::mpsc::channel();
4501        let interface = Arc::new(MockInterface::new(tx));
4502        let session_manager = Arc::new(SessionManager::new(interface.clone(), 512));
4503        let block_server = BlockServer::new(512, session_manager.clone());
4504
4505        let sm_completer = session_manager.clone();
4506        std::thread::spawn(move || {
4507            while let Ok(req) = rx.recv() {
4508                if let Operation::Read { vmo_offset, block_count, .. } = req.operation {
4509                    if let Some(vmo) = req.vmo {
4510                        let data = vec![0xABu8; (block_count * 512) as usize];
4511                        vmo.write(&data, vmo_offset).unwrap();
4512                    }
4513                }
4514                sm_completer.complete_request(req.request_id, Ok(()));
4515            }
4516        });
4517
4518        let (mapper_proxy, mapper_stream) =
4519            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
4520        let scope = fasync::Scope::new();
4521        scope.spawn(async move {
4522            let _ = block_server.handle_mapper_requests(mapper_stream).await;
4523        });
4524
4525        // 1. Prepare and open intermediate root session without pager.
4526        let (root_session_proxy, root_session_server) =
4527            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4528        let root_mapping_vmo = zx::Vmo::create(65536).unwrap();
4529
4530        // Register partition in root session at key 500:
4531        // Logical 0..8192 maps to physical device offset 4096..12288.
4532        let partition_key = 500u64;
4533        let extents =
4534            mapping::Extents::try_new([mapping::Extent::new(0..8192, Some(4096))], 0).unwrap();
4535        let partition_extents: Vec<u64> = mapping::Extents::encode_extents(&extents).collect();
4536        let mut payload_bytes = Vec::new();
4537        for w in &partition_extents {
4538            payload_bytes.extend_from_slice(&w.to_le_bytes());
4539        }
4540        let cmd = mapping::RawMappingCommand {
4541            opcode: mapping::MAPPINGS_COMMAND,
4542            offset: 0,
4543            key: partition_key,
4544            stored_size: 8192,
4545            device_offset: 0,
4546            metadata_count: 0,
4547            extent_count: partition_extents.len() as u32,
4548        };
4549        let mut sender = vmo_fifo::SyncSender::<mapping::RawMappingCommand>::new(
4550            root_mapping_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4551            1024,
4552            256,
4553        )
4554        .unwrap();
4555        let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
4556        payload_buf.data().copy_from_slice(&payload_bytes);
4557        payload_buf.commit(cmd).unwrap();
4558
4559        let res = mapper_proxy
4560            .open_session(root_session_server, root_mapping_vmo, None, None)
4561            .await
4562            .unwrap();
4563        assert_matches!(res, Ok(()));
4564
4565        // 2. Open child session with pager through partition key 500.
4566        let pager = Arc::new(zx::Pager::create(zx::PagerOptions::empty()).unwrap());
4567        let port = zx::Port::create();
4568        let child_key = 2002u64;
4569        let paged_vmo = pager.create_vmo(zx::VmoOptions::empty(), &port, child_key, 4096).unwrap();
4570
4571        let child_mapping_vmo = zx::Vmo::create(65536).unwrap();
4572        let delivery_queue = zx::Vmo::create(65536).unwrap();
4573        let vmo_provider = Arc::new(blob_pager_and_verifier::TestVmoProvider::new(
4574            pager.clone(),
4575            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4576        ));
4577        vmo_provider
4578            .register_vmo(child_key, paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap());
4579        let receiver = vmo_fifo::Receiver::<mapping::RawDeliveryCommand>::new(
4580            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4581            mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
4582        )
4583        .unwrap();
4584        let _delivery_processor = blob_pager_and_verifier::DeliveryQueueProcessor::spawn(
4585            receiver,
4586            vmo_provider,
4587            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4588        )
4589        .unwrap();
4590
4591        let extents =
4592            mapping::Extents::try_new([mapping::Extent::new(0..4096, Some(0))], 0).unwrap();
4593        let child_extents: Vec<u64> = mapping::Extents::encode_extents(&extents).collect();
4594        let mut child_payload_bytes = Vec::new();
4595        for w in &child_extents {
4596            child_payload_bytes.extend_from_slice(&w.to_le_bytes());
4597        }
4598        let child_cmd = mapping::RawMappingCommand {
4599            opcode: mapping::MAPPINGS_COMMAND,
4600            offset: 0,
4601            key: child_key,
4602            stored_size: 4096,
4603            device_offset: 0,
4604            metadata_count: 0,
4605            extent_count: child_extents.len() as u32,
4606        };
4607        let mut child_sender = vmo_fifo::SyncSender::<mapping::RawMappingCommand>::new(
4608            child_mapping_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4609            1024,
4610            256,
4611        )
4612        .unwrap();
4613        let mut child_payload_buf =
4614            child_sender.reserve_payload(child_payload_bytes.len()).unwrap();
4615        child_payload_buf.data().copy_from_slice(&child_payload_bytes);
4616        child_payload_buf.commit(child_cmd).unwrap();
4617
4618        let mut child_res = Err(zx::Status::NOT_FOUND);
4619        for _ in 0..100 {
4620            let (_child_session_proxy, child_session_server) =
4621                fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4622            let child_res_raw = root_session_proxy
4623                .open_child_session(
4624                    child_session_server,
4625                    child_mapping_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4626                    partition_key,
4627                    Some(port.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap()),
4628                    Some(delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap()),
4629                )
4630                .await
4631                .unwrap();
4632            if child_res_raw.is_ok() {
4633                child_res = Ok(());
4634                break;
4635            }
4636            fasync::Timer::new(std::time::Duration::from_millis(10)).await;
4637        }
4638        assert_matches!(child_res, Ok(()));
4639
4640        let reader_thread = std::thread::spawn(move || {
4641            let mut buf = [0u8; 4096];
4642            paged_vmo.read(&mut buf, 0).expect("paged vmo read failed");
4643            buf
4644        });
4645
4646        let read_bytes = reader_thread.join().unwrap();
4647        assert_eq!(read_bytes.len(), 4096);
4648        assert_eq!(read_bytes, [0xABu8; 4096]);
4649    }
4650
4651    #[fuchsia::test]
4652    async fn test_mapper_child_session_closed_when_parent_closed() {
4653        use crate::callback_interface::SessionManager;
4654        use crate::testing::MockInterface;
4655        use fidl::endpoints::Proxy as _;
4656
4657        let (tx, _rx) = std::sync::mpsc::channel();
4658        let interface = Arc::new(MockInterface::new(tx));
4659        let session_manager = Arc::new(SessionManager::new(interface.clone(), 512));
4660        let block_server = BlockServer::new(512, session_manager.clone());
4661
4662        let (mapper_proxy, mapper_stream) =
4663            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
4664        let scope = fasync::Scope::new();
4665        scope.spawn(async move {
4666            let _ = block_server.handle_mapper_requests(mapper_stream).await;
4667        });
4668
4669        // 1. Prepare and open intermediate root session without pager.
4670        let (root_session_proxy, root_session_server) =
4671            fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4672        let root_mapping_vmo = zx::Vmo::create(65536).unwrap();
4673
4674        let partition_key = 500u64;
4675        let extents =
4676            mapping::Extents::try_new([mapping::Extent::new(0..8192, Some(4096))], 0).unwrap();
4677        let partition_extents: Vec<u64> = mapping::Extents::encode_extents(&extents).collect();
4678        let mut payload_bytes = Vec::new();
4679        for w in &partition_extents {
4680            payload_bytes.extend_from_slice(&w.to_le_bytes());
4681        }
4682        let cmd = mapping::RawMappingCommand {
4683            opcode: mapping::MAPPINGS_COMMAND,
4684            offset: 0,
4685            key: partition_key,
4686            stored_size: 8192,
4687            device_offset: 0,
4688            metadata_count: 0,
4689            extent_count: partition_extents.len() as u32,
4690        };
4691        let mut sender = vmo_fifo::SyncSender::<mapping::RawMappingCommand>::new(
4692            root_mapping_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4693            1024,
4694            256,
4695        )
4696        .unwrap();
4697        let mut payload_buf = sender.reserve_payload(payload_bytes.len()).unwrap();
4698        payload_buf.data().copy_from_slice(&payload_bytes);
4699        payload_buf.commit(cmd).unwrap();
4700
4701        let res = mapper_proxy
4702            .open_session(root_session_server, root_mapping_vmo, None, None)
4703            .await
4704            .unwrap();
4705        assert_matches!(res, Ok(()));
4706
4707        // 2. Open child session.
4708        let child_mapping_vmo = zx::Vmo::create(65536).unwrap();
4709        let mut child_session = None;
4710        for _ in 0..100 {
4711            let (child_session_proxy, child_session_server) =
4712                fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4713            let child_res_raw = root_session_proxy
4714                .open_child_session(
4715                    child_session_server,
4716                    child_mapping_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
4717                    partition_key,
4718                    None,
4719                    None,
4720                )
4721                .await
4722                .unwrap();
4723            if child_res_raw.is_ok() {
4724                child_session = Some(child_session_proxy);
4725                break;
4726            }
4727            fasync::Timer::new(std::time::Duration::from_millis(10)).await;
4728        }
4729        let child_session_proxy = child_session.expect("Failed to open child session");
4730
4731        // Close the parent session.
4732        root_session_proxy.close().await.unwrap().expect("Close parent session failed");
4733
4734        // The child session should be closed as well.
4735        child_session_proxy.on_closed().await.unwrap();
4736    }
4737
4738    #[fuchsia::test]
4739    async fn test_mapper_session_detaches_vmo_on_close() {
4740        let interface = Arc::new(MockInterface::default());
4741        let block_server = BlockServer::new(BLOCK_SIZE, interface.clone());
4742
4743        let (mapper_proxy, mapper_stream) =
4744            fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
4745        let server_fut = async move {
4746            block_server.handle_mapper_requests(mapper_stream).await.unwrap();
4747        };
4748
4749        let interface_ref = &interface;
4750        let client_fut = async move {
4751            let (session_proxy, session_server) =
4752                fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
4753            let mapping_vmo = zx::Vmo::create(4096).unwrap();
4754            mapper_proxy
4755                .open_session(session_server, mapping_vmo, None, None)
4756                .await
4757                .unwrap()
4758                .unwrap();
4759
4760            let _paged_vmo = session_proxy
4761                .create_vmo(1, 4096, fblock::CreateVmoOptions::empty())
4762                .await
4763                .unwrap()
4764                .unwrap();
4765            assert_eq!(interface_ref.attached_vmos.load(Ordering::Relaxed), 1);
4766
4767            session_proxy.close().await.unwrap().unwrap();
4768        };
4769
4770        futures::join!(server_fut, client_fut);
4771        assert_eq!(interface.attached_vmos.load(Ordering::Relaxed), 0);
4772    }
4773}