1mod delivery;
40
41pub use delivery::{
42 DELIVERY_DATA_SIZE, DeliveryQueueProcessor, DeliveryQueueProvider, TestVmoProvider,
43 UnverifiedPages,
44};
45
46use anyhow::{Context, Error, anyhow, bail};
47use event_listener as _;
48use fidl_fuchsia_storage_block as fblock;
49use fidl_fuchsia_storage_mapping as fmapping;
50use fuchsia_async as fasync;
51use fuchsia_hash::{HASH_SIZE, Hash};
52use fuchsia_merkle::{MerkleVerifier, ReadSizedMerkleVerifier};
53use fuchsia_sync::Mutex;
54use std::collections::{HashMap, hash_map};
55use std::ops::Range;
56use std::sync::{Arc, OnceLock, Weak};
57use storage_ptr_slice::PtrByteSlice;
58use zx;
59
60struct ZeroChildrenReceiver {
63 blob: OnceLock<Weak<CachedBlob>>,
64}
65
66impl fasync::PacketReceiver for ZeroChildrenReceiver {
67 fn receive_packet(&self, packet: zx::Packet) {
68 if let zx::PacketContents::SignalOne(signals) = packet.contents() {
69 if signals.observed().contains(zx::Signals::VMO_ZERO_CHILDREN) {
70 if let Some(cached_blob) = self.blob.get().and_then(|w| w.upgrade()) {
71 if let Some(cache) = cached_blob.cache.upgrade() {
72 cache.on_zero_children(&cached_blob);
73 }
74 }
75 }
76 }
77 }
78}
79
80struct CachedBlob {
88 vmo: zx::Vmo,
91 root_hash: [u8; 32],
93 cache: Weak<Cache>,
96 vmo_key: u64,
99 len: u64,
101 registration: fasync::ReceiverRegistration<ZeroChildrenReceiver>,
103 merkle_verifier: OnceLock<ReadSizedMerkleVerifier>,
105}
106
107impl CachedBlob {
108 fn create_child(&self) -> Result<zx::Vmo, zx::Status> {
109 self.vmo.create_child(zx::VmoChildOptions::REFERENCE | zx::VmoChildOptions::NO_WRITE, 0, 0)
110 }
111
112 fn wait_for_zero_children(&self) {
113 let _ = self.vmo.wait_async(
114 fasync::EHandle::local().port(),
115 self.registration.key(),
116 zx::Signals::VMO_ZERO_CHILDREN,
117 zx::WaitAsyncOpts::empty(),
118 );
119 }
120}
121
122impl Drop for CachedBlob {
123 fn drop(&mut self) {
124 if let Some(cache) = self.cache.upgrade() {
125 let session = cache.mapping_session.clone();
126 let vmo_key = self.vmo_key;
127 cache.scope.spawn(async move {
128 let _ = session.close(vmo_key).await;
129 });
130 }
131 }
132}
133
134enum BlobState {
135 Opening {
139 root_hash: [u8; 32],
140 event: Arc<event_listener::Event>,
141 verifier: Option<ReadSizedMerkleVerifier>,
142 },
143 Ready(Arc<CachedBlob>),
146}
147
148struct CacheState {
149 next_key: u64,
150 keys_by_hash: HashMap<[u8; 32], u64>,
152 blob_states_by_key: HashMap<u64, BlobState>,
153}
154
155impl Default for CacheState {
156 fn default() -> Self {
157 Self {
158 next_key: 1,
159 keys_by_hash: HashMap::default(),
160 blob_states_by_key: HashMap::default(),
161 }
162 }
163}
164
165impl CacheState {
166 fn generate_key(&mut self) -> u64 {
167 let key = self.next_key;
168 self.next_key += 1;
169 key
170 }
171}
172
173struct PendingEntry {
179 cache: Arc<Cache>,
180 key: u64,
181}
182
183impl PendingEntry {
184 fn key(&self) -> u64 {
185 self.key
186 }
187
188 fn commit(self, vmo: zx::Vmo, size: u64) -> Arc<CachedBlob> {
191 let zero_children_receiver = ZeroChildrenReceiver { blob: OnceLock::new() };
192 let zero_children_registration =
193 fasync::EHandle::local().register_receiver(zero_children_receiver);
194
195 let (cached, event) = {
196 let mut state = self.cache.state.lock();
197 let (root_hash, event, verifier) = match state.blob_states_by_key.remove(&self.key) {
198 Some(BlobState::Opening { root_hash, event, verifier }) => {
199 (root_hash, event, verifier)
200 }
201 _ => {
202 unreachable!("Blob for key {} was unexpectedly not in Opening state", self.key)
203 }
204 };
205
206 let merkle_verifier = OnceLock::new();
207 if let Some(verifier) = verifier {
208 let _ = merkle_verifier.set(verifier);
209 }
210
211 let cached = Arc::new(CachedBlob {
212 vmo,
213 root_hash,
214 cache: Arc::downgrade(&self.cache),
215 vmo_key: self.key,
216 len: size,
217 registration: zero_children_registration,
218 merkle_verifier,
219 });
220
221 cached
222 .registration
223 .receiver()
224 .blob
225 .set(Arc::downgrade(&cached))
226 .expect("Failed to bind CachedBlob to ZeroChildrenReceiver: already initialized");
227
228 state.blob_states_by_key.insert(self.key, BlobState::Ready(cached.clone()));
229
230 (cached, event)
231 };
232
233 cached.wait_for_zero_children();
234 event.notify(usize::MAX);
235 cached
236 }
237}
238
239impl Drop for PendingEntry {
240 fn drop(&mut self) {
241 let event = {
242 let mut state = self.cache.state.lock();
243 match state.blob_states_by_key.get(&self.key) {
244 Some(BlobState::Ready(_)) => return,
246 Some(BlobState::Opening { .. }) => {
247 if let Some(BlobState::Opening { root_hash, event, .. }) =
248 state.blob_states_by_key.remove(&self.key)
249 {
250 state.keys_by_hash.remove(&root_hash);
251 event
252 } else {
253 return;
254 }
255 }
256 None => {
257 log::error!(
258 key = self.key;
259 "PendingEntry::drop: key unexpectedly absent from cache"
260 );
261 return;
262 }
263 }
264 };
265
266 let session = self.cache.mapping_session.clone();
272 let key = self.key;
273 self.cache.scope.spawn(async move {
274 let _ = session.close(key).await;
275 });
276
277 event.notify(usize::MAX);
279 }
280}
281
282enum CacheLookup {
284 Ready(zx::Vmo),
286 Pending(event_listener::EventListener),
289 Missing(PendingEntry),
292}
293
294struct Cache {
296 state: Mutex<CacheState>,
297 mapping_session: fmapping::MappingSessionProxy,
298 pager: Arc<zx::Pager>,
299 delivery_vmo: zx::Vmo,
300 scope: fasync::ScopeHandle,
301}
302
303impl Cache {
304 fn new(
305 mapping_session: fmapping::MappingSessionProxy,
306 pager: Arc<zx::Pager>,
307 delivery_vmo: zx::Vmo,
308 scope: fasync::ScopeHandle,
309 ) -> Self {
310 Self {
311 state: Mutex::new(CacheState::default()),
312 mapping_session,
313 pager,
314 delivery_vmo,
315 scope,
316 }
317 }
318
319 fn delivery_vmo(&self) -> &zx::Vmo {
320 &self.delivery_vmo
321 }
322
323 fn get_by_key(&self, key: u64) -> Option<Arc<CachedBlob>> {
324 let state = self.state.lock();
325 match state.blob_states_by_key.get(&key) {
326 Some(BlobState::Ready(cached)) => Some(cached.clone()),
327 _ => None,
328 }
329 }
330
331 #[cfg(test)]
332 fn is_merkle_initialized(&self, key: u64) -> bool {
333 let state = self.state.lock();
334 match state.blob_states_by_key.get(&key) {
335 Some(BlobState::Opening { verifier, .. }) => verifier.is_some(),
336 Some(BlobState::Ready(blob)) => blob.merkle_verifier.get().is_some(),
337 None => false,
338 }
339 }
340
341 fn lookup(self: &Arc<Self>, identifier: &[u8; 32]) -> Result<CacheLookup, Error> {
349 let mut state = self.state.lock();
350 if let Some(&key) = state.keys_by_hash.get(identifier) {
351 match state.blob_states_by_key.get(&key) {
352 Some(BlobState::Ready(cached)) => {
353 let child = cached
354 .create_child()
355 .map_err(|s| anyhow!("Failed to create child VMO: {s}"))?;
356 return Ok(CacheLookup::Ready(child));
357 }
358 Some(BlobState::Opening { event, .. }) => {
359 return Ok(CacheLookup::Pending(event.listen()));
360 }
361 None => unreachable!("keys_by_hash pointed to non-existent key {key}"),
362 }
363 }
364
365 let event = Arc::new(event_listener::Event::new());
366 let key = state.generate_key();
367 state.keys_by_hash.insert(*identifier, key);
368 state
369 .blob_states_by_key
370 .insert(key, BlobState::Opening { root_hash: *identifier, event, verifier: None });
371 Ok(CacheLookup::Missing(PendingEntry { cache: self.clone(), key }))
372 }
373
374 fn on_zero_children(&self, blob: &Arc<CachedBlob>) {
375 let mut state = self.state.lock();
380 if let Ok(info) = blob.vmo.info() {
381 if info.num_children > 0 {
382 blob.wait_for_zero_children();
385 return;
386 }
387 }
388 if let hash_map::Entry::Occupied(entry) = state.keys_by_hash.entry(blob.root_hash) {
390 if *entry.get() == blob.vmo_key {
391 entry.remove();
392 }
393 }
394 if let hash_map::Entry::Occupied(entry) = state.blob_states_by_key.entry(blob.vmo_key) {
396 if let BlobState::Ready(cached) = entry.get() {
397 if Arc::ptr_eq(cached, blob) {
398 entry.remove();
399 }
400 }
401 }
402 }
403
404 fn create_merkle_verifier(
405 root_hash: &[u8; 32],
406 vmo_key: u64,
407 hashes: Box<[Hash]>,
408 ) -> Result<ReadSizedMerkleVerifier, Error> {
409 let hashes = if hashes.is_empty() { Box::new([Hash::from(*root_hash)]) } else { hashes };
414
415 match MerkleVerifier::new(Hash::from(*root_hash), hashes) {
416 Ok(verifier) => ReadSizedMerkleVerifier::new(verifier, delivery::DELIVERY_DATA_SIZE)
417 .map_err(|e| anyhow!("Failed to create ReadSizedMerkleVerifier: {e:?}")),
418 Err(e) => {
419 bail!("Failed to verify merkle leaves for key {vmo_key}: {e:?}");
420 }
421 }
422 }
423
424 fn report_pager_failure(&self, vmo: &zx::Vmo, range: Range<u64>, status: zx::Status) {
425 let pager_status = match status {
428 zx::Status::IO_DATA_INTEGRITY => zx::Status::IO_DATA_INTEGRITY,
429 zx::Status::NO_SPACE => zx::Status::NO_SPACE,
430 zx::Status::BUFFER_TOO_SMALL | zx::Status::FILE_BIG => zx::Status::BUFFER_TOO_SMALL,
431 zx::Status::IO
432 | zx::Status::IO_DATA_LOSS
433 | zx::Status::IO_INVALID
434 | zx::Status::IO_MISSED_DEADLINE
435 | zx::Status::IO_NOT_PRESENT
436 | zx::Status::IO_OVERRUN
437 | zx::Status::IO_REFUSED
438 | zx::Status::PEER_CLOSED => zx::Status::IO,
439 _ => zx::Status::BAD_STATE,
440 };
441
442 if let Err(error) = self.pager.op_range(zx::PagerOp::Fail(pager_status), vmo, range) {
445 log::error!(error:?; "Failed to report pager failure to kernel");
446 }
447 }
448}
449
450impl delivery::DeliveryQueueProvider for Cache {
451 fn deliver_pages(
452 &self,
453 key: u64,
454 target_offset: u64,
455 unverified_pages: UnverifiedPages<'_>,
456 ) -> Result<(), Error> {
457 let cached_blob = self
458 .get_by_key(key)
459 .ok_or_else(|| anyhow!("Unknown or expired key {key} in deliver_pages"))?;
460
461 let page_size = zx::system_get_page_size() as u64;
462 let page_aligned_size = cached_blob.len.div_ceil(page_size) * page_size;
463
464 let verifier = match cached_blob.merkle_verifier.get() {
465 Some(v) => v,
466 None => {
467 self.report_pager_failure(
472 &cached_blob.vmo,
473 0..page_aligned_size,
474 zx::Status::BAD_STATE,
475 );
476 bail!("Data received for uninitialized blob {}", key);
477 }
478 };
479
480 let chunk_len = unverified_pages.len_in_bytes() as u64;
481
482 let bounded_len = std::cmp::min(chunk_len, page_aligned_size.saturating_sub(target_offset));
487 let chunk_range = target_offset..target_offset + bounded_len;
488
489 let remaining_in_blob = cached_blob.len.saturating_sub(target_offset);
493 let unaligned_len = std::cmp::min(chunk_len, remaining_in_blob) as usize;
494 if let Err(e) = verifier.verify_aligned(
495 target_offset as usize,
496 unverified_pages.as_ptr_byte_slice(),
497 unaligned_len,
498 ) {
499 self.report_pager_failure(&cached_blob.vmo, chunk_range, zx::Status::IO_DATA_INTEGRITY);
500 bail!("Failed to verify payload for blob {}: {:?}", key, e);
501 }
502
503 if let Err(e) = self.pager.supply_pages(
504 &cached_blob.vmo,
505 chunk_range.clone(),
506 &self.delivery_vmo,
507 unverified_pages.delivery_offset(),
508 ) {
509 self.report_pager_failure(&cached_blob.vmo, chunk_range, e);
510 bail!("Failed to supply pages: {:?}", e);
511 }
512
513 Ok(())
514 }
515
516 fn register_blob(&self, key: u64, leaf_data: PtrByteSlice<'_>) -> Result<(), Error> {
517 if leaf_data.len() % HASH_SIZE != 0 {
518 bail!("RegisterBlob invalid leaf length must be a multiple of HASH_SIZE");
519 }
520
521 let hashes: Vec<Hash> = (0..(leaf_data.len() / HASH_SIZE))
522 .map(|i| {
523 let chunk = leaf_data.subslice(i * HASH_SIZE..(i + 1) * HASH_SIZE);
524 let mut hash_bytes = [0u8; HASH_SIZE];
525 chunk.copy_to_slice(&mut hash_bytes);
526 Hash::from(hash_bytes)
527 })
528 .collect();
529
530 let mut state = self.state.lock();
531 match state.blob_states_by_key.get_mut(&key) {
532 Some(BlobState::Opening { root_hash, verifier, .. }) => {
533 if verifier.is_some() {
534 log::warn!(key;
535 "RegisterBlob received for blob but it is already initialized"
536 );
537 return Ok(());
538 }
539 *verifier =
540 Some(Self::create_merkle_verifier(root_hash, key, hashes.into_boxed_slice())?);
541 Ok(())
542 }
543 Some(BlobState::Ready(cached)) => {
544 if cached.merkle_verifier.get().is_some() {
545 log::warn!(key;
546 "RegisterBlob received for blob but it is already initialized"
547 );
548 return Ok(());
549 }
550 let verifier = Self::create_merkle_verifier(
551 &cached.root_hash,
552 key,
553 hashes.into_boxed_slice(),
554 )?;
555 let _ = cached.merkle_verifier.set(verifier);
556 Ok(())
557 }
558 None => {
559 bail!("RegisterBlob received for unknown or expired key {key}");
560 }
561 }
562 }
563}
564
565pub struct BlobPagerAndVerifier {
573 _delivery_processor: delivery::DeliveryQueueProcessor,
575 _mapper_session: fblock::MapperSessionProxy,
577 port: zx::Port,
578 pager: Arc<zx::Pager>,
579 cache: Arc<Cache>,
580 _scope: fasync::Scope,
581}
582
583impl BlobPagerAndVerifier {
584 pub fn delivery_vmo(&self) -> &zx::Vmo {
586 self.cache.delivery_vmo()
587 }
588
589 pub async fn new(
597 mapping_provider: &fmapping::MappingProviderProxy,
598 mapper: &fblock::MapperProxy,
599 ) -> Result<Self, Error> {
600 let (mapping_session, session_server) =
602 fidl::endpoints::create_proxy::<fmapping::MappingSessionMarker>();
603 let mapping_vmo = mapping_provider
604 .open_session(session_server)
605 .await
606 .context("FIDL error calling MappingProvider.OpenSession")?
607 .map_err(|e| anyhow!("MappingProvider open_session failed: {e:?}"))?;
608
609 let port = zx::Port::create();
611 let delivery_queue = zx::Vmo::create(mapping::DELIVERY_VMO_SIZE)
612 .context("Failed to create delivery queue VMO")?;
613 let delivery_vmo: zx::Vmo = delivery_queue
614 .duplicate_handle(zx::Rights::SAME_RIGHTS)
615 .context("Failed to duplicate delivery queue VMO")?;
616 let delivery_queue_dup = delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS)?;
617 let receiver = vmo_fifo::Receiver::<mapping::RawDeliveryCommand>::new(
618 delivery_queue,
619 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
620 )
621 .context("Failed to create delivery queue receiver")?;
622 let pager = Arc::new(
623 zx::Pager::create(zx::PagerOptions::empty()).context("Failed to create pager")?,
624 );
625 let port_dup = port.duplicate_handle(zx::Rights::SAME_RIGHTS)?;
626 let (mapper_session, mapper_session_server) =
627 fidl::endpoints::create_proxy::<fblock::MapperSessionMarker>();
628 mapper
629 .open_session(
630 mapper_session_server,
631 mapping_vmo,
632 Some(port_dup),
633 Some(delivery_queue_dup),
634 )
635 .await
636 .context("FIDL error calling Mapper.OpenSession")?
637 .map_err(|e| anyhow!("Mapper.OpenSession failed: {e:?}"))?;
638
639 let scope = fasync::Scope::new();
640 let cache = Arc::new(Cache::new(
641 mapping_session,
642 pager.clone(),
643 delivery_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)?,
644 scope.to_handle(),
645 ));
646 let _delivery_processor =
647 delivery::DeliveryQueueProcessor::spawn(receiver, cache.clone(), delivery_vmo)?;
648
649 Ok(Self {
650 _delivery_processor,
651 _mapper_session: mapper_session,
652 port,
653 pager,
654 cache,
655 _scope: scope,
656 })
657 }
658
659 pub async fn create_vmo(&self, identifier: &[u8; 32]) -> Result<zx::Vmo, Error> {
661 loop {
662 let entry = match self.cache.lookup(identifier)? {
663 CacheLookup::Ready(child) => return Ok(child),
664 CacheLookup::Pending(listener) => {
665 listener.await;
666 continue; }
668 CacheLookup::Missing(entry) => entry,
669 };
670
671 let key = entry.key();
672
673 let size = self
676 .cache
677 .mapping_session
678 .open(key, identifier)
679 .await
680 .context("FIDL error calling MappingSession.Open")?
681 .map_err(|e| anyhow!("MappingSession.Open failed for blob: {e:?}"))?;
682
683 let vmo = self.pager.create_vmo(zx::VmoOptions::empty(), &self.port, key, size)?;
684
685 let first_child = vmo
688 .create_child(zx::VmoChildOptions::REFERENCE | zx::VmoChildOptions::NO_WRITE, 0, 0)
689 .map_err(|s| anyhow!("Failed to create child VMO: {}", s))?;
690
691 entry.commit(vmo, size);
692
693 return Ok(first_child);
694 }
695 }
696}
697
698#[cfg(test)]
699mod tests {
700 use super::*;
701 use crate::delivery::DELIVERY_DATA_SIZE;
702 use futures::TryStreamExt;
703 use futures::channel::oneshot;
704 use mapping::{DeliveryCommand, RawDeliveryCommand};
705 use vmo_fifo::SyncSender;
706
707 const TEST_BLOB_SIZE: u64 = (DELIVERY_DATA_SIZE * 2) as u64;
708 const TEST_VMO_KEY: u64 = 1;
709
710 struct TestEnv {
715 pager_and_verifier: Arc<BlobPagerAndVerifier>,
716 mapping_task: fasync::Task<()>,
717 mapper_task: fasync::Task<()>,
718 close_signal: Option<oneshot::Receiver<()>>,
719 delivery_vmo: zx::Vmo,
720 pub valid_root: [u8; 32],
721 pub valid_leaves: Vec<u8>,
722 pub blob_data: Vec<u8>,
723 }
724
725 impl TestEnv {
726 async fn new(blob_size: u64) -> Self {
727 let blob_data = vec![0x42u8; blob_size as usize];
728 let (root, leaf_hashes) =
729 fuchsia_merkle::MerkleRootBuilder::new(Vec::new()).complete(&blob_data);
730 let expected_hash: [u8; 32] = root.into();
731
732 let mut flat_leaves = Vec::new();
733 for hash in &leaf_hashes {
734 flat_leaves.extend_from_slice(hash.as_bytes());
735 }
736
737 let (mapping_proxy, mut mapping_stream) =
738 fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
739
740 let (close_tx, close_rx) = oneshot::channel();
741 let mapping_task = fasync::Task::spawn(async move {
742 let mut close_tx = Some(close_tx);
743 if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
744 mapping_stream.try_next().await.expect("try_next failed")
745 {
746 let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
747 .expect("zx::Vmo::create failed");
748 responder.send(Ok(mapping_vmo)).expect("send failed");
749
750 let mut session_stream = session.into_stream();
751 let mut open_calls = 0;
752 while let Some(request) =
753 session_stream.try_next().await.expect("try_next failed")
754 {
755 match request {
756 fmapping::MappingSessionRequest::Open {
757 key,
758 identifier,
759 responder,
760 } => {
761 assert_eq!(identifier, expected_hash);
762 assert_eq!(key, TEST_VMO_KEY);
763 open_calls += 1;
764 assert_eq!(
765 open_calls, 1,
766 "Open should only be called once per identifier"
767 );
768 responder.send(Ok(blob_size)).expect("send failed");
769 }
770 fmapping::MappingSessionRequest::Close { key, responder } => {
771 assert_eq!(key, TEST_VMO_KEY);
772 responder.send(Ok(())).expect("send failed");
773 if let Some(tx) = close_tx.take() {
774 let _ = tx.send(());
775 }
776 }
777 _ => {}
778 }
779 }
780 }
781 });
782
783 let (mapper_proxy, mut mapper_stream) =
784 fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
785
786 let (tx_delivery, rx_delivery) = oneshot::channel();
787 let mapper_task = fasync::Task::spawn(async move {
788 if let Some(fblock::MapperRequest::OpenSession {
789 delivery_queue, responder, ..
790 }) = mapper_stream.try_next().await.expect("try_next failed")
791 {
792 tx_delivery.send(delivery_queue).expect("send failed");
793 responder.send(Ok(())).expect("send failed");
794 }
795 });
796
797 Self {
798 pager_and_verifier: Arc::new(
799 BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
800 .await
801 .expect("BlobPagerAndVerifier::new failed"),
802 ),
803 mapping_task,
804 mapper_task,
805 close_signal: Some(close_rx),
806 delivery_vmo: rx_delivery
807 .await
808 .expect("rx_delivery wait failed")
809 .expect("delivery_vmo was None"),
810 valid_root: expected_hash,
811 valid_leaves: flat_leaves,
812 blob_data,
813 }
814 }
815
816 async fn teardown(self) {
817 drop(self.pager_and_verifier);
818 self.mapping_task.await;
819 self.mapper_task.await;
820 }
821 }
822
823 struct ExpectedReadThread {
824 started_rx: Option<oneshot::Receiver<()>>,
825 done_rx: oneshot::Receiver<()>,
826 handle: std::thread::JoinHandle<()>,
827 }
828
829 impl ExpectedReadThread {
830 async fn wait_started(&mut self) {
831 if let Some(rx) = self.started_rx.take() {
832 rx.await.expect("reader thread failed to start");
833 }
834 }
835
836 async fn wait_and_verify(self) {
837 self.done_rx.await.expect("reader thread panicked or hung");
838 self.handle.join().expect("join failed");
839 }
840 }
841
842 fn spawn_reader_expect_error(
843 vmo: &zx::Vmo,
844 offset: u64,
845 size: usize,
846 expected_status: zx::Status,
847 ) -> ExpectedReadThread {
848 let vmo_clone =
849 vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate_handle failed");
850 let (started_tx, started_rx) = oneshot::channel();
851 let (done_tx, done_rx) = oneshot::channel();
852
853 let handle = std::thread::spawn(move || {
854 let mut buf = vec![0u8; size];
855 let _ = started_tx.send(());
856 let err = vmo_clone.read(&mut buf, offset).expect_err("read should fail");
857 assert_eq!(err, expected_status);
858 let _ = done_tx.send(());
859 });
860
861 ExpectedReadThread { started_rx: Some(started_rx), done_rx, handle }
862 }
863
864 #[fuchsia::test]
865 async fn test_create_vmo() {
866 let env = TestEnv::new(TEST_BLOB_SIZE).await;
867 let vmo =
868 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("Failed to create VMO");
869 let size = vmo.get_size().expect("get_size failed");
870
871 let page_size = zx::system_get_page_size() as u64;
873 let expected_pages = (TEST_BLOB_SIZE + page_size - 1) / page_size;
874 assert_eq!(size, expected_pages * page_size);
875
876 drop(vmo);
880
881 env.teardown().await;
882 }
883
884 #[fuchsia::test]
885 async fn test_create_vmo_concurrent_access() {
886 let env = TestEnv::new(TEST_BLOB_SIZE).await;
887 let mut futures = vec![];
888 for _ in 0..10 {
889 let verifier = env.pager_and_verifier.clone();
890 futures.push(fasync::Task::spawn(async move {
891 verifier.create_vmo(&env.valid_root).await.expect("Failed to create VMO")
892 }));
893 }
894 let vmos = futures::future::join_all(futures).await;
895
896 drop(vmos);
897 env.teardown().await;
898 }
899
900 #[fuchsia::test]
901 async fn test_zero_children_eviction() {
902 let mut env = TestEnv::new(TEST_BLOB_SIZE).await;
903
904 let child_vmo =
905 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
906
907 {
908 let state = env.pager_and_verifier.cache.state.lock();
909 assert_eq!(state.keys_by_hash.get(&env.valid_root), Some(&TEST_VMO_KEY));
910 assert!(matches!(
911 state.blob_states_by_key.get(&TEST_VMO_KEY),
912 Some(BlobState::Ready(_))
913 ));
914 }
915
916 drop(child_vmo);
917
918 env.close_signal.take().expect("Missing signal").await.expect("Failed to close");
921
922 {
923 let state = env.pager_and_verifier.cache.state.lock();
924 assert!(state.keys_by_hash.get(&env.valid_root).is_none());
925 assert!(state.blob_states_by_key.get(&TEST_VMO_KEY).is_none());
926 }
927
928 env.teardown().await;
929 }
930
931 #[fuchsia::test]
932 async fn test_create_vmo_multiple_distinct_blobs() {
933 use futures::StreamExt;
934
935 let (mapping_proxy, mut mapping_stream) =
936 fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
937 let (mapper_proxy, mut mapper_stream) =
938 fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
939
940 let (close_tx, mut close_rx) = futures::channel::mpsc::unbounded();
941 let mapping_task = fasync::Task::spawn(async move {
942 if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
943 mapping_stream.try_next().await.expect("try_next failed")
944 {
945 let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
946 .expect("zx::Vmo::create failed");
947 responder.send(Ok(mapping_vmo)).expect("send failed");
948
949 let mut session_stream = session.into_stream();
950 while let Some(request) = session_stream.try_next().await.expect("try_next failed")
951 {
952 match request {
953 fmapping::MappingSessionRequest::Open { responder, .. } => {
954 responder.send(Ok(8192)).expect("send failed");
955 }
956 fmapping::MappingSessionRequest::Close { key, responder } => {
957 responder.send(Ok(())).expect("send failed");
958 let _ = close_tx.unbounded_send(key);
959 }
960 _ => {}
961 }
962 }
963 }
964 });
965
966 let mapper_task = fasync::Task::spawn(async move {
967 if let Some(fblock::MapperRequest::OpenSession { responder, .. }) =
968 mapper_stream.try_next().await.expect("try_next failed")
969 {
970 responder.send(Ok(())).expect("send failed");
971 }
972 });
973
974 let pager_and_verifier = Arc::new(
975 BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
976 .await
977 .expect("BlobPagerAndVerifier::new failed"),
978 );
979
980 let hash1: [u8; 32] = [0x11; 32];
981 let hash2: [u8; 32] = [0x22; 32];
982 let hash3: [u8; 32] = [0x33; 32];
983
984 let vmo1 = pager_and_verifier.create_vmo(&hash1).await.expect("create_vmo 1 failed");
985 let vmo2 = pager_and_verifier.create_vmo(&hash2).await.expect("create_vmo 2 failed");
986 let vmo3 = pager_and_verifier.create_vmo(&hash3).await.expect("create_vmo 3 failed");
987
988 let (key1, key2, key3) = {
989 let state = pager_and_verifier.cache.state.lock();
990 let k1 = *state.keys_by_hash.get(&hash1).expect("hash1 missing");
991 let k2 = *state.keys_by_hash.get(&hash2).expect("hash2 missing");
992 let k3 = *state.keys_by_hash.get(&hash3).expect("hash3 missing");
993 assert_ne!(k1, k2);
994 assert_ne!(k2, k3);
995 assert_ne!(k1, k3);
996 assert!(matches!(state.blob_states_by_key.get(&k1), Some(BlobState::Ready(_))));
997 assert!(matches!(state.blob_states_by_key.get(&k2), Some(BlobState::Ready(_))));
998 assert!(matches!(state.blob_states_by_key.get(&k3), Some(BlobState::Ready(_))));
999 (k1, k2, k3)
1000 };
1001
1002 drop(vmo2);
1004 let closed_key = close_rx.next().await.expect("expected close signal");
1005 assert_eq!(closed_key, key2);
1006
1007 {
1009 let state = pager_and_verifier.cache.state.lock();
1010 assert!(state.keys_by_hash.get(&hash2).is_none());
1011 assert!(state.blob_states_by_key.get(&key2).is_none());
1012 assert_eq!(state.keys_by_hash.get(&hash1), Some(&key1));
1013 assert_eq!(state.keys_by_hash.get(&hash3), Some(&key3));
1014 assert!(matches!(state.blob_states_by_key.get(&key1), Some(BlobState::Ready(_))));
1015 assert!(matches!(state.blob_states_by_key.get(&key3), Some(BlobState::Ready(_))));
1016 }
1017
1018 let vmo2_new = pager_and_verifier.create_vmo(&hash2).await.expect("re-create vmo 2 failed");
1020 {
1021 let state = pager_and_verifier.cache.state.lock();
1022 let k = *state.keys_by_hash.get(&hash2).expect("hash2 missing after re-open");
1023 assert_ne!(k, key2);
1024 assert!(matches!(state.blob_states_by_key.get(&k), Some(BlobState::Ready(_))));
1025 }
1026
1027 drop(vmo1);
1028 drop(vmo3);
1029 drop(vmo2_new);
1030
1031 drop(pager_and_verifier);
1032 mapping_task.await;
1033 mapper_task.await;
1034 }
1035
1036 #[fuchsia::test]
1037 async fn test_cache_is_cleared_when_create_vmo_fails() {
1038 let (mapping_proxy, mut mapping_stream) =
1039 fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
1040 let (mapper_proxy, mut mapper_stream) =
1041 fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
1042
1043 let (mock_opened_tx, mock_opened_rx) = oneshot::channel();
1044 let mapping_task = fasync::Task::spawn(async move {
1045 let mut mock_opened_tx = Some(mock_opened_tx);
1046 if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
1047 mapping_stream.try_next().await.expect("try_next failed")
1048 {
1049 let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
1050 .expect("zx::Vmo::create failed");
1051 responder.send(Ok(mapping_vmo)).expect("send failed");
1052 let mut session_stream = session.into_stream();
1053
1054 if let Some(fmapping::MappingSessionRequest::Open {
1055 responder: _responder, ..
1056 }) = session_stream.try_next().await.expect("try_next failed")
1057 {
1058 if let Some(tx) = mock_opened_tx.take() {
1060 let _ = tx.send(());
1061 }
1062
1063 let () = std::future::pending().await;
1065 }
1066 }
1067 });
1068
1069 let mapper_task = fasync::Task::spawn(async move {
1070 if let Some(fblock::MapperRequest::OpenSession { responder, .. }) =
1071 mapper_stream.try_next().await.expect("try_next failed")
1072 {
1073 responder.send(Ok(())).expect("send failed");
1074 }
1075 });
1076
1077 let pager_and_verifer = Arc::new(
1078 BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
1079 .await
1080 .expect("BlobPagerAndVerifier::new failed"),
1081 );
1082
1083 let blob_data = vec![0x42u8; TEST_BLOB_SIZE as usize];
1084 let (root, _) = fuchsia_merkle::MerkleRootBuilder::new(Vec::new()).complete(&blob_data);
1085 let hash_val: [u8; 32] = root.into();
1086
1087 let verifier_clone = pager_and_verifer.clone();
1088
1089 let identifier = hash_val;
1090 let (abortable_future, abort_handle) = futures::future::abortable(async move {
1091 let _ = verifier_clone.create_vmo(&identifier).await;
1092 });
1093
1094 let identifier = &hash_val;
1095 let primary_task = fasync::Task::spawn(abortable_future);
1096
1097 mock_opened_rx.await.expect("Open exited abruptly");
1100 {
1101 let state = pager_and_verifer.cache.state.lock();
1102 assert_eq!(state.keys_by_hash.get(identifier), Some(&1));
1103 assert!(matches!(state.blob_states_by_key.get(&1), Some(BlobState::Opening { .. })));
1104 }
1105
1106 abort_handle.abort();
1108 let _ = primary_task.await;
1109
1110 {
1112 let state = pager_and_verifer.cache.state.lock();
1113 assert!(state.keys_by_hash.get(identifier).is_none());
1114 assert!(state.blob_states_by_key.get(&1).is_none());
1115 }
1116
1117 drop(pager_and_verifer);
1118 drop(mapping_task);
1119 drop(mapper_task);
1120 }
1121
1122 #[fuchsia::test]
1123 async fn test_create_vmo_error_clears_cache_and_unblocks_waiters() {
1124 let (mapping_proxy, mut mapping_stream) =
1125 fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
1126 let (mapper_proxy, mut mapper_stream) =
1127 fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
1128
1129 let (open_received_tx, open_received_rx) = oneshot::channel();
1130 let (reply_tx, reply_rx) = oneshot::channel();
1131
1132 let mapping_task = fasync::Task::spawn(async move {
1133 if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
1134 mapping_stream.try_next().await.expect("try_next failed")
1135 {
1136 let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
1137 .expect("zx::Vmo::create failed");
1138 responder.send(Ok(mapping_vmo)).expect("send failed");
1139 let mut session_stream = session.into_stream();
1140
1141 if let Some(fmapping::MappingSessionRequest::Open { responder, .. }) =
1142 session_stream.try_next().await.expect("try_next failed")
1143 {
1144 open_received_tx.send(()).expect("send open_received failed");
1146
1147 reply_rx.await.expect("reply_rx failed");
1149 responder.send(Err(zx::Status::NOT_FOUND.into_raw())).expect("send failed");
1150 }
1151 }
1153 });
1154
1155 let mapper_task = fasync::Task::spawn(async move {
1156 if let Some(fblock::MapperRequest::OpenSession { responder, .. }) =
1157 mapper_stream.try_next().await.expect("try_next failed")
1158 {
1159 responder.send(Ok(())).expect("send failed");
1160 }
1161 });
1162
1163 let pager_and_verifier = Arc::new(
1164 BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
1165 .await
1166 .expect("BlobPagerAndVerifier::new failed"),
1167 );
1168
1169 let hash: [u8; 32] = [0x55; 32];
1170
1171 let verifier1 = pager_and_verifier.clone();
1173 let primary_future = fasync::Task::spawn(async move { verifier1.create_vmo(&hash).await });
1174
1175 open_received_rx.await.expect("open_received_rx failed");
1177 {
1178 let state = pager_and_verifier.cache.state.lock();
1179 assert!(matches!(
1180 state.keys_by_hash.get(&hash).and_then(|k| state.blob_states_by_key.get(k)),
1181 Some(BlobState::Opening { .. })
1182 ));
1183 }
1184
1185 let verifier2 = pager_and_verifier.clone();
1187 let secondary_future =
1188 fasync::Task::spawn(async move { verifier2.create_vmo(&hash).await });
1189
1190 reply_tx.send(()).expect("reply_tx send failed");
1192
1193 assert!(primary_future.await.is_err());
1195 assert!(secondary_future.await.is_err());
1196
1197 {
1199 let state = pager_and_verifier.cache.state.lock();
1200 assert!(state.keys_by_hash.get(&hash).is_none());
1201 assert!(state.blob_states_by_key.is_empty());
1202 }
1203
1204 drop(pager_and_verifier);
1205 drop(mapping_task);
1206 drop(mapper_task);
1207 }
1208
1209 #[fuchsia::test]
1210 async fn test_register_blob_arrives_before_open_completes() {
1211 let (mapping_proxy, mut mapping_stream) =
1215 fidl::endpoints::create_proxy_and_stream::<fmapping::MappingProviderMarker>();
1216 let (mapper_proxy, mut mapper_stream) =
1217 fidl::endpoints::create_proxy_and_stream::<fblock::MapperMarker>();
1218
1219 let blob_data = vec![0x42u8; TEST_BLOB_SIZE as usize];
1220 let (root, leaf_hashes) =
1221 fuchsia_merkle::MerkleRootBuilder::new(Vec::new()).complete(&blob_data);
1222 let expected_hash: [u8; 32] = root.into();
1223
1224 let mut flat_leaves = Vec::new();
1225 for hash in &leaf_hashes {
1226 flat_leaves.extend_from_slice(hash.as_bytes());
1227 }
1228
1229 let (tx_delivery, rx_delivery) = oneshot::channel();
1230 let (open_key_tx, open_key_rx) = oneshot::channel();
1231 let (resume_open_tx, resume_open_rx) = oneshot::channel();
1232 let (close_tx, close_rx) = oneshot::channel();
1233
1234 let mapping_task = fasync::Task::spawn(async move {
1235 let mut close_tx = Some(close_tx);
1236 let mut open_key_tx = Some(open_key_tx);
1237 let mut resume_open_rx = Some(resume_open_rx);
1238 if let Some(fmapping::MappingProviderRequest::OpenSession { session, responder }) =
1239 mapping_stream.try_next().await.expect("try_next failed")
1240 {
1241 let mapping_vmo = zx::Vmo::create(zx::system_get_page_size().into())
1242 .expect("zx::Vmo::create failed");
1243 responder.send(Ok(mapping_vmo)).expect("send failed");
1244
1245 let mut session_stream = session.into_stream();
1246 while let Some(request) = session_stream.try_next().await.expect("try_next failed")
1247 {
1248 match request {
1249 fmapping::MappingSessionRequest::Open { key, identifier, responder } => {
1250 assert_eq!(identifier, expected_hash);
1251 if let Some(tx) = open_key_tx.take() {
1253 tx.send(key).expect("open_key_tx failed");
1254 }
1255 if let Some(rx) = resume_open_rx.take() {
1256 rx.await.expect("resume_open_rx failed");
1257 }
1258 responder.send(Ok(TEST_BLOB_SIZE)).expect("send failed");
1259 }
1260 fmapping::MappingSessionRequest::Close { key: _, responder } => {
1261 responder.send(Ok(())).expect("send failed");
1262 if let Some(tx) = close_tx.take() {
1263 let _ = tx.send(());
1264 }
1265 }
1266 _ => {}
1267 }
1268 }
1269 }
1270 });
1271
1272 let mapper_task = fasync::Task::spawn(async move {
1273 if let Some(fblock::MapperRequest::OpenSession { delivery_queue, responder, .. }) =
1274 mapper_stream.try_next().await.expect("try_next failed")
1275 {
1276 tx_delivery.send(delivery_queue).expect("send failed");
1277 responder.send(Ok(())).expect("send failed");
1278 }
1279 });
1280
1281 let pager_and_verifier = Arc::new(
1282 BlobPagerAndVerifier::new(&mapping_proxy, &mapper_proxy)
1283 .await
1284 .expect("BlobPagerAndVerifier::new failed"),
1285 );
1286
1287 let delivery_vmo =
1288 rx_delivery.await.expect("rx_delivery wait failed").expect("delivery_vmo was None");
1289 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1290 delivery_vmo
1291 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1292 .expect("duplicate_handle failed"),
1293 zx::system_get_page_size() as usize,
1294 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1295 )
1296 .expect("SyncSender::new failed");
1297
1298 let verifier_clone = pager_and_verifier.clone();
1299 let create_task =
1300 fasync::Task::spawn(async move { verifier_clone.create_vmo(&expected_hash).await });
1301
1302 let key = open_key_rx.await.expect("open_key_rx failed");
1304
1305 let mut payload =
1307 sender.reserve_payload(flat_leaves.len()).expect("reserve_payload failed");
1308 payload.data().copy_from_slice(&flat_leaves);
1309 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1310 key: key as u64,
1311 offset: payload.offset(),
1312 length: flat_leaves.len() as u32,
1313 }
1314 .into();
1315 payload.commit(raw_cmd).expect("commit failed");
1316
1317 while !pager_and_verifier.cache.is_merkle_initialized(key) {
1319 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1320 }
1321
1322 resume_open_tx.send(()).expect("resume_open_tx failed");
1324 let vmo = create_task.await.expect("create_vmo failed");
1325
1326 let chunk = &blob_data[..DELIVERY_DATA_SIZE];
1328 let mut data_payload = sender.reserve_payload(chunk.len()).expect("reserve_payload failed");
1329 data_payload.data().copy_from_slice(chunk);
1330 let data_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1331 key: key as u64,
1332 offset: data_payload.offset(),
1333 length: chunk.len() as u32,
1334 target_offset: 0,
1335 }
1336 .into();
1337 data_payload.commit(data_cmd).expect("commit failed");
1338
1339 let mut read_buf = vec![0u8; DELIVERY_DATA_SIZE];
1340 vmo.read(&mut read_buf, 0).expect("read failed");
1341 assert_eq!(read_buf, chunk);
1342
1343 drop(vmo);
1344 close_rx.await.expect("close_rx failed");
1345
1346 drop(pager_and_verifier);
1347 mapping_task.await;
1348 mapper_task.await;
1349 }
1350
1351 #[fuchsia::test]
1352 async fn test_delivery_register_blob_verification() {
1353 let env = TestEnv::new(TEST_BLOB_SIZE).await;
1354
1355 let _paged_vmo =
1356 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1357
1358 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1359 env.delivery_vmo
1360 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1361 .expect("duplicate_handle failed"),
1362 zx::system_get_page_size() as usize, mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1364 )
1365 .expect("SyncSender::new failed");
1366
1367 let mut corrupted_leaves = env.valid_leaves.clone();
1369 corrupted_leaves[0] ^= 0xFF;
1370
1371 let mut bad_payload =
1372 sender.reserve_payload(corrupted_leaves.len()).expect("reserve_payload failed");
1373 bad_payload.data().copy_from_slice(&corrupted_leaves);
1374
1375 let bad_raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1376 key: TEST_VMO_KEY,
1377 offset: bad_payload.offset(),
1378 length: corrupted_leaves.len() as u32,
1379 }
1380 .into();
1381 bad_payload.commit(bad_raw_cmd).expect("commit failed");
1382
1383 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1385
1386 let blob =
1387 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1388
1389 assert!(blob.merkle_verifier.get().is_none());
1391
1392 let mut payload =
1394 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1395 payload.data().copy_from_slice(&env.valid_leaves);
1396
1397 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1398 key: TEST_VMO_KEY,
1399 offset: payload.offset(),
1400 length: env.valid_leaves.len() as u32,
1401 }
1402 .into();
1403 payload.commit(raw_cmd).expect("commit failed");
1404
1405 while blob.merkle_verifier.get().is_none() {
1406 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1407 }
1408
1409 drop(_paged_vmo);
1410 env.teardown().await;
1411 }
1412
1413 #[fuchsia::test]
1414 async fn test_register_blob_invalid_commands() {
1415 let env = TestEnv::new(TEST_BLOB_SIZE).await;
1416
1417 let _paged_vmo =
1418 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1419
1420 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1421 env.delivery_vmo
1422 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1423 .expect("duplicate_handle failed"),
1424 std::mem::align_of::<RawDeliveryCommand>(),
1425 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1426 )
1427 .expect("SyncSender::new failed");
1428
1429 let blob =
1430 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1431
1432 let mut invalid_key_payload =
1434 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1435 invalid_key_payload.data().copy_from_slice(&env.valid_leaves);
1436 let invalid_key_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1437 key: 9999 as u64,
1438 offset: invalid_key_payload.offset(),
1439 length: env.valid_leaves.len() as u32,
1440 }
1441 .into();
1442 invalid_key_payload.commit(invalid_key_cmd).expect("commit failed");
1443 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1444 assert!(blob.merkle_verifier.get().is_none());
1445
1446 let oob_payload = sender.reserve_payload(8).expect("reserve_payload failed");
1448 let oob_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1449 key: TEST_VMO_KEY,
1450 offset: u32::MAX - 4, length: 8,
1452 }
1453 .into();
1454 oob_payload.commit(oob_cmd).expect("commit failed");
1455 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1456 assert!(blob.merkle_verifier.get().is_none());
1457
1458 let invalid_len_payload = sender.reserve_payload(10).expect("reserve_payload failed");
1460 let invalid_len_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1461 key: TEST_VMO_KEY,
1462 offset: invalid_len_payload.offset(),
1463 length: 10,
1464 }
1465 .into();
1466 invalid_len_payload.commit(invalid_len_cmd).expect("commit failed");
1467 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1468 assert!(blob.merkle_verifier.get().is_none());
1469
1470 drop(_paged_vmo);
1471 env.teardown().await;
1472 }
1473
1474 #[fuchsia::test]
1475 async fn test_delivery_data_supplies_pages() {
1476 let env = TestEnv::new(TEST_BLOB_SIZE).await;
1477
1478 let paged_vmo =
1479 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1480
1481 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1482 env.delivery_vmo
1483 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1484 .expect("duplicate_handle failed"),
1485 zx::system_get_page_size() as usize, mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1487 )
1488 .expect("SyncSender::new failed");
1489
1490 let mut payload =
1491 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1492 payload.data().copy_from_slice(&env.valid_leaves);
1493 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1494 key: TEST_VMO_KEY,
1495 offset: payload.offset(),
1496 length: env.valid_leaves.len() as u32,
1497 }
1498 .into();
1499 payload.commit(raw_cmd).expect("commit failed");
1500
1501 let blob =
1502 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1503 while blob.merkle_verifier.get().is_none() {
1504 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1505 }
1506
1507 let test_payload = &env.blob_data[..DELIVERY_DATA_SIZE];
1508
1509 let mut payload =
1510 sender.reserve_payload(test_payload.len()).expect("reserve_payload failed");
1511 payload.data().copy_from_slice(&test_payload);
1512
1513 let cmd: RawDeliveryCommand = DeliveryCommand::Data {
1514 key: TEST_VMO_KEY,
1515 target_offset: 0,
1516 length: test_payload.len() as u32,
1517 offset: payload.offset(),
1518 }
1519 .into();
1520 payload.commit(cmd).expect("commit failed");
1521
1522 let mut read_buf = vec![0u8; DELIVERY_DATA_SIZE];
1523 paged_vmo.read(&mut read_buf, 0).expect("read paged_vmo failed");
1524 assert_eq!(read_buf, test_payload);
1525
1526 drop(paged_vmo);
1527 env.teardown().await;
1528 }
1529
1530 #[fuchsia::test]
1531 async fn test_delivery_data_corrupted() {
1532 let blob_size = DELIVERY_DATA_SIZE as u64;
1535 let env = TestEnv::new(blob_size).await;
1536
1537 let paged_vmo =
1538 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1539
1540 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1541 env.delivery_vmo
1542 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1543 .expect("duplicate_handle failed"),
1544 zx::system_get_page_size() as usize,
1545 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1546 )
1547 .expect("SyncSender::new failed");
1548
1549 let mut payload =
1550 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1551 payload.data().copy_from_slice(&env.valid_leaves);
1552 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1553 key: TEST_VMO_KEY,
1554 offset: payload.offset(),
1555 length: env.valid_leaves.len() as u32,
1556 }
1557 .into();
1558 payload.commit(raw_cmd).expect("commit failed");
1559
1560 let blob =
1561 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1562 while blob.merkle_verifier.get().is_none() {
1563 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1564 }
1565
1566 let mut corrupted_data = env.blob_data.clone();
1568 corrupted_data[0] ^= 0xFF;
1569
1570 let mut payload =
1571 sender.reserve_payload(corrupted_data.len()).expect("reserve_payload failed");
1572 payload.data().copy_from_slice(&corrupted_data);
1573 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1574 key: TEST_VMO_KEY,
1575 offset: payload.offset(),
1576 length: corrupted_data.len() as u32,
1577 target_offset: 0,
1578 }
1579 .into();
1580
1581 let page_size = zx::system_get_page_size() as usize;
1584 let mut reader1 =
1585 spawn_reader_expect_error(&paged_vmo, 0, page_size, zx::Status::IO_DATA_INTEGRITY);
1586 let mut reader2 = spawn_reader_expect_error(
1587 &paged_vmo,
1588 page_size as u64,
1589 page_size,
1590 zx::Status::IO_DATA_INTEGRITY,
1591 );
1592
1593 reader1.wait_started().await;
1596 reader2.wait_started().await;
1597 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1598
1599 payload.commit(raw_cmd).expect("commit failed");
1602
1603 reader1.wait_and_verify().await;
1604 reader2.wait_and_verify().await;
1605
1606 drop(paged_vmo);
1607 env.teardown().await;
1608 }
1609
1610 #[fuchsia::test]
1611 async fn test_delivery_data_corrupted_chunk_preserves_supplied_pages() {
1612 let chunk_size = DELIVERY_DATA_SIZE;
1613 let blob_size = (chunk_size * 2) as u64;
1614 let env = TestEnv::new(blob_size).await;
1615
1616 let paged_vmo =
1617 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1618
1619 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1620 env.delivery_vmo
1621 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1622 .expect("duplicate_handle failed"),
1623 zx::system_get_page_size() as usize,
1624 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1625 )
1626 .expect("SyncSender::new failed");
1627
1628 let mut payload =
1629 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1630 payload.data().copy_from_slice(&env.valid_leaves);
1631 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1632 key: TEST_VMO_KEY,
1633 offset: payload.offset(),
1634 length: env.valid_leaves.len() as u32,
1635 }
1636 .into();
1637 payload.commit(raw_cmd).expect("commit failed");
1638
1639 let blob =
1640 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1641 while blob.merkle_verifier.get().is_none() {
1642 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1643 }
1644
1645 let page_size = zx::system_get_page_size() as usize;
1646
1647 let paged_vmo_clone1 =
1650 paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate_handle failed");
1651 let (tx1, rx1) = oneshot::channel();
1652 let expected_page0 = env.blob_data[..page_size].to_vec();
1653 let thread1 = std::thread::spawn(move || {
1654 let mut buf = vec![0u8; page_size];
1655 paged_vmo_clone1.read(&mut buf, 0).expect("read should succeed");
1656 assert_eq!(buf, expected_page0);
1657 let _ = tx1.send(());
1658 });
1659 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1660
1661 let valid_chunk = &env.blob_data[..chunk_size];
1662 let mut payload =
1663 sender.reserve_payload(valid_chunk.len()).expect("reserve_payload failed");
1664 payload.data().copy_from_slice(valid_chunk);
1665 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1666 key: TEST_VMO_KEY,
1667 offset: payload.offset(),
1668 length: valid_chunk.len() as u32,
1669 target_offset: 0,
1670 }
1671 .into();
1672 payload.commit(raw_cmd).expect("commit failed");
1673
1674 rx1.await.expect("thread1 panicked or hung");
1675 thread1.join().expect("join failed");
1676
1677 let mut reader2 = spawn_reader_expect_error(
1679 &paged_vmo,
1680 chunk_size as u64,
1681 page_size,
1682 zx::Status::IO_DATA_INTEGRITY,
1683 );
1684 reader2.wait_started().await;
1685 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1686
1687 let mut corrupted_chunk = env.blob_data[chunk_size..].to_vec();
1689 corrupted_chunk[0] ^= 0xFF;
1690
1691 let mut payload =
1692 sender.reserve_payload(corrupted_chunk.len()).expect("reserve_payload failed");
1693 payload.data().copy_from_slice(&corrupted_chunk);
1694 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1695 key: TEST_VMO_KEY,
1696 offset: payload.offset(),
1697 length: corrupted_chunk.len() as u32,
1698 target_offset: chunk_size as u64,
1699 }
1700 .into();
1701 payload.commit(raw_cmd).expect("commit failed");
1702
1703 reader2.wait_and_verify().await;
1704
1705 let mut buf = vec![0u8; page_size];
1707 paged_vmo.read(&mut buf, 0).expect("previously supplied page 0 must still be readable");
1708 assert_eq!(buf, env.blob_data[..page_size]);
1709
1710 drop(paged_vmo);
1711 env.teardown().await;
1712 }
1713
1714 #[fuchsia::test]
1715 async fn test_delivery_data_corrupted_multi_chunk_read() {
1716 let chunk_size = DELIVERY_DATA_SIZE;
1717 let blob_size = (chunk_size * 3) as u64;
1718 let env = TestEnv::new(blob_size).await;
1719
1720 let paged_vmo =
1721 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1722
1723 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1724 env.delivery_vmo
1725 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1726 .expect("duplicate_handle failed"),
1727 zx::system_get_page_size() as usize,
1728 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1729 )
1730 .expect("SyncSender::new failed");
1731
1732 let mut payload =
1733 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1734 payload.data().copy_from_slice(&env.valid_leaves);
1735 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1736 key: TEST_VMO_KEY,
1737 offset: payload.offset(),
1738 length: env.valid_leaves.len() as u32,
1739 }
1740 .into();
1741 payload.commit(raw_cmd).expect("commit failed");
1742
1743 let blob =
1744 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1745 while blob.merkle_verifier.get().is_none() {
1746 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1747 }
1748
1749 let mut reader = spawn_reader_expect_error(
1755 &paged_vmo,
1756 0,
1757 blob_size as usize,
1758 zx::Status::IO_DATA_INTEGRITY,
1759 );
1760 reader.wait_started().await;
1761 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1762
1763 let valid_chunk1 = &env.blob_data[..chunk_size];
1765 let mut payload =
1766 sender.reserve_payload(valid_chunk1.len()).expect("reserve_payload failed");
1767 payload.data().copy_from_slice(valid_chunk1);
1768 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1769 key: TEST_VMO_KEY,
1770 offset: payload.offset(),
1771 length: valid_chunk1.len() as u32,
1772 target_offset: 0,
1773 }
1774 .into();
1775 payload.commit(raw_cmd).expect("commit failed");
1776
1777 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1780
1781 let mut corrupted_chunk2 = env.blob_data[chunk_size..chunk_size * 2].to_vec();
1783 corrupted_chunk2[0] ^= 0xFF;
1784 let mut payload =
1785 sender.reserve_payload(corrupted_chunk2.len()).expect("reserve_payload failed");
1786 payload.data().copy_from_slice(&corrupted_chunk2);
1787 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1788 key: TEST_VMO_KEY,
1789 offset: payload.offset(),
1790 length: corrupted_chunk2.len() as u32,
1791 target_offset: chunk_size as u64,
1792 }
1793 .into();
1794 payload.commit(raw_cmd).expect("commit failed");
1795
1796 reader.wait_and_verify().await;
1800
1801 drop(paged_vmo);
1802 env.teardown().await;
1803 }
1804
1805 #[fuchsia::test]
1806 async fn test_delivery_data_uninitialized() {
1807 let env = TestEnv::new(TEST_BLOB_SIZE).await;
1808
1809 let paged_vmo =
1810 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1811
1812 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1813 env.delivery_vmo
1814 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1815 .expect("duplicate_handle failed"),
1816 zx::system_get_page_size() as usize,
1817 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1818 )
1819 .expect("SyncSender::new failed");
1820
1821 let mut reader = spawn_reader_expect_error(&paged_vmo, 0, 8192, zx::Status::BAD_STATE);
1823 reader.wait_started().await;
1824 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1825
1826 let chunk = &env.blob_data[..8192];
1828 let mut payload = sender.reserve_payload(chunk.len()).expect("reserve_payload failed");
1829 payload.data().copy_from_slice(chunk);
1830 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1831 key: TEST_VMO_KEY,
1832 offset: payload.offset(),
1833 length: chunk.len() as u32,
1834 target_offset: 0,
1835 }
1836 .into();
1837 payload.commit(raw_cmd).expect("commit failed");
1838
1839 reader.wait_and_verify().await;
1840
1841 drop(paged_vmo);
1842 env.teardown().await;
1843 }
1844
1845 #[fuchsia::test]
1846 async fn test_delivery_data_multiple_chunks() {
1847 let env = TestEnv::new(TEST_BLOB_SIZE).await;
1848
1849 let paged_vmo =
1850 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("create_vmo failed");
1851
1852 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1853 env.delivery_vmo
1854 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1855 .expect("duplicate_handle failed"),
1856 zx::system_get_page_size() as usize,
1857 mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1858 )
1859 .expect("SyncSender::new failed");
1860
1861 let mut payload =
1862 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1863 payload.data().copy_from_slice(&env.valid_leaves);
1864 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1865 key: TEST_VMO_KEY,
1866 offset: payload.offset(),
1867 length: env.valid_leaves.len() as u32,
1868 }
1869 .into();
1870 payload.commit(raw_cmd).expect("commit failed");
1871
1872 let blob =
1873 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1874 while blob.merkle_verifier.get().is_none() {
1875 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1876 }
1877
1878 let expected_data = env.blob_data.clone();
1879 let paged_vmo_clone =
1880 paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).expect("duplicate_handle failed");
1881
1882 let (tx, rx) = oneshot::channel();
1883 let thread = std::thread::spawn(move || {
1884 let mut buf = vec![0u8; expected_data.len()];
1885 paged_vmo_clone.read(&mut buf, 0).expect("failed to read from paged vmo");
1886 assert_eq!(buf, expected_data);
1887 let _ = tx.send(());
1888 });
1889
1890 let chunk_size = DELIVERY_DATA_SIZE;
1892 for (i, chunk) in env.blob_data.chunks(chunk_size).enumerate() {
1893 let mut payload = sender.reserve_payload(chunk.len()).expect("reserve_payload failed");
1894 payload.data().copy_from_slice(chunk);
1895 let raw_cmd: RawDeliveryCommand = DeliveryCommand::Data {
1896 key: TEST_VMO_KEY,
1897 offset: payload.offset(),
1898 length: chunk.len() as u32,
1899 target_offset: (i * chunk_size) as u64,
1900 }
1901 .into();
1902 payload.commit(raw_cmd).expect("commit failed");
1903 }
1904
1905 rx.await.expect("reading thread panicked or hung");
1906 thread.join().expect("join failed");
1907
1908 drop(paged_vmo);
1909 env.teardown().await;
1910 }
1911
1912 #[fuchsia::test]
1913 async fn test_delivery_data_unaligned_blob_size() {
1914 let data_size = (DELIVERY_DATA_SIZE + 1024) as u64; let env = TestEnv::new(data_size).await;
1916
1917 let paged_vmo =
1918 env.pager_and_verifier.create_vmo(&env.valid_root).await.expect("Failed to create VMO");
1919
1920 let mut sender = SyncSender::<RawDeliveryCommand>::new(
1921 env.delivery_vmo
1922 .duplicate_handle(zx::Rights::SAME_RIGHTS)
1923 .expect("duplicate_handle failed"),
1924 zx::system_get_page_size() as usize, mapping::PENDING_DELIVERY_COMMANDS_CAPACITY,
1926 )
1927 .expect("SyncSender::new failed");
1928
1929 let mut payload =
1930 sender.reserve_payload(env.valid_leaves.len()).expect("reserve_payload failed");
1931 payload.data().copy_from_slice(&env.valid_leaves);
1932 let raw_cmd: RawDeliveryCommand = DeliveryCommand::RegisterBlob {
1933 key: TEST_VMO_KEY,
1934 offset: payload.offset(),
1935 length: env.valid_leaves.len() as u32,
1936 }
1937 .into();
1938 payload.commit(raw_cmd).expect("commit failed");
1939
1940 let blob =
1941 env.pager_and_verifier.cache.get_by_key(TEST_VMO_KEY).expect("get_by_key failed");
1942 while blob.merkle_verifier.get().is_none() {
1943 fasync::Timer::new(std::time::Duration::from_millis(5)).await;
1944 }
1945
1946 let test_payload = env.blob_data.clone();
1947 assert_eq!(test_payload.len(), data_size as usize);
1948
1949 let chunk1 = &test_payload[..DELIVERY_DATA_SIZE];
1950 let mut payload1 = sender.reserve_payload(chunk1.len()).expect("reserve_payload failed");
1951 payload1.data().copy_from_slice(chunk1);
1952 let off1 = payload1.offset();
1953 payload1
1954 .commit(
1955 DeliveryCommand::Data {
1956 key: TEST_VMO_KEY,
1957 target_offset: 0,
1958 length: chunk1.len() as u32,
1959 offset: off1,
1960 }
1961 .into(),
1962 )
1963 .expect("commit failed");
1964
1965 let chunk2 = &test_payload[DELIVERY_DATA_SIZE..];
1966 let page_size = zx::system_get_page_size() as usize;
1967 let chunk2_len_aligned = chunk2.len().div_ceil(page_size) * page_size;
1968 let mut payload2 =
1969 sender.reserve_payload(chunk2_len_aligned).expect("reserve_payload failed");
1970 let payload2_data = payload2.data();
1971 payload2_data.subslice_mut(0..chunk2.len()).copy_from_slice(chunk2);
1972 payload2_data.subslice_mut(chunk2.len()..chunk2_len_aligned).fill(0);
1973
1974 let off2 = payload2.offset();
1975 payload2
1976 .commit(
1977 DeliveryCommand::Data {
1978 key: TEST_VMO_KEY,
1979 target_offset: DELIVERY_DATA_SIZE as u64,
1980 length: chunk2_len_aligned as u32,
1981 offset: off2,
1982 }
1983 .into(),
1984 )
1985 .expect("commit failed");
1986
1987 let mut read_buf = vec![0u8; data_size as usize];
1988 paged_vmo.read(&mut read_buf, 0).expect("read paged_vmo failed");
1989 assert_eq!(read_buf, test_payload);
1990
1991 let page_size = zx::system_get_page_size() as u64;
1992 let out_of_bounds_offset = data_size.div_ceil(page_size) * page_size;
1994 let mut read_buf = vec![0u8; 1];
1996 let res = paged_vmo.read(&mut read_buf, out_of_bounds_offset);
1997 assert_eq!(res, Err(zx::Status::OUT_OF_RANGE));
1998
1999 drop(paged_vmo);
2000 env.teardown().await;
2001 }
2002}