1use 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 Block(BlockInfo),
41 Partition(PartitionInfo),
43 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 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#[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#[derive(Clone, Default, Debug)]
118pub struct PartitionInfo {
119 pub device_flags: fblock::DeviceFlag,
121 pub max_transfer_blocks: Option<NonZero<u32>>,
122 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 pub flags: Option<u64>,
132}
133
134#[derive(Clone, Default, Debug)]
136pub struct VolumeInfo {
137 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
146struct 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 compressed_range: Range<usize>,
162
163 uncompressed_range: Range<u64>,
165
166 bytes_so_far: u64,
167 mapping: Arc<VmoMapping>,
168 buffer: Option<Buffer<'static>>,
169}
170
171impl DecompressionInfo {
172 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 server_ep: zx::EventPair,
186}
187
188#[derive(Default)]
191pub struct KeyRegistry {
192 keys: Mutex<HashMap<zx::Koid, RegisteredHwKey>>,
193}
194
195impl KeyRegistry {
196 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 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 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
250impl<S> ActiveRequestsInner<S> {
252 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 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 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 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 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 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
357pub 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
370pub 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#[derive(Clone, Debug, Default, PartialEq, Eq)]
388pub struct OffsetMap {
389 mappings: Vec<BlockOffsetMapping>,
390}
391
392impl OffsetMap {
393 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 pub fn empty() -> Self {
410 Self { mappings: Vec::new() }
411 }
412
413 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 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
515pub trait SessionManager: 'static {
518 type Orchestrator: Borrow<Self> + Send + Sync;
523
524 const SUPPORTS_DECOMPRESSION: bool;
525
526 type Session;
527
528 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 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 fn get_info(&self) -> Cow<'_, DeviceInfo>;
555
556 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 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 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 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 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 fn active_requests(&self) -> &ActiveRequests<Self::Session>;
604
605 fn key_registry(&self) -> &KeyRegistry;
607}
608
609pub 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 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 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 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 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 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 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#[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 unsafe {
994 let _ = fuchsia_runtime::vmar_root_self().unmap(self.base, self.size);
995 }
996 }
997}
998
999enum HandleRequestResult {
1000 Ok,
1002 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 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 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 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 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 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 if group_or_request.is_group() {
1244 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 if group.status.is_ok() {
1254 group.status = Err(zx::Status::INVALID_ARGS);
1255 }
1256 return Err(None);
1258 }
1259 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 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 let Err(s) = group.status {
1319 operation = Err(s);
1320 }
1321 } else if group.status.is_err() {
1322 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 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 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 registered_vmo
1377 .get_or_create_mapping()
1378 .and_then(|mapping| {
1379 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 operation = Ok(Operation::StartDecompressedRead {
1396 required_buffer_size,
1397 device_block_offset,
1398 block_count,
1399 options,
1400 });
1401 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 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 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 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 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 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
1564pub 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 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 CloseVmo,
1598 StartDecompressedRead {
1600 required_buffer_size: usize,
1601 device_block_offset: u64,
1602 block_count: u32,
1603 options: ReadOptions,
1604 },
1605 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 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 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 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 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 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 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 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 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 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 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 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 std::mem::drop(proxy);
2207
2208 session_proxy.close().await.unwrap().unwrap();
2209
2210 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 *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 *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 *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 *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 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 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 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 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 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 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 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, 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 poll_fn(|cx: &mut Context<'_>| {
2812 if waiting_readers.lock().len() == 2 {
2813 Poll::Ready(())
2814 } else {
2815 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 start_rx.next().await.unwrap();
2905
2906 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_tx.send(()).unwrap();
2918
2919 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 poll_fn(|cx: &mut Context<'_>| {
3052 if counter.load(Ordering::Relaxed) == 1 {
3053 Poll::Ready(())
3054 } else {
3055 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 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 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 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 tx.send(()).unwrap();
3203
3204 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 assert!(fasync::TestExecutor::poll_until_stalled(&mut fut).await.is_pending());
3279
3280 let mut fut2 = pin!(proxy.get_volume_info());
3281
3282 assert!(fasync::TestExecutor::poll_until_stalled(&mut fut2).await.is_pending());
3284
3285 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 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 event.1.store(true, Ordering::SeqCst);
3372 event.0.notify(usize::MAX);
3373 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 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 *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 *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 *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 *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 *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 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 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 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 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 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 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 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 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 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 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 #[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(), 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 #[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(), 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 #[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 Arc::new(MockInterface::default()),
4017 );
4018 block_server.handle_requests(stream).await.unwrap();
4019 },
4020 async move {
4021 async fn settle() {
4022 for _ in 0..32 {
4024 fasync::yield_now().await;
4025 }
4026 }
4027
4028 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 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 {
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 {
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 {
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 {
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 {
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 {
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 {
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 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 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 assert_eq!(map.map(10), Some((110, 70)));
4255 assert_eq!(map.map(55), Some((155, 25)));
4258 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 {
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 {
4293 let (session_proxy, server) = fidl::endpoints::create_proxy();
4294 let mapping = fblock::BlockOffsetMapping {
4295 target_block_offset: u64::MAX - 10, 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 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 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 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 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 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 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 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 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 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 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 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 root_session_proxy.close().await.unwrap().expect("Close parent session failed");
4833
4834 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 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 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 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 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 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}