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