1#![cfg(target_os = "fuchsia")]
6#![deny(missing_docs)]
7
8use crate::types::{BlobId, BlobInfo};
11use fidl_fuchsia_fxfs as ffxfs;
12use fidl_fuchsia_pkg as fpkg;
13use fuchsia_pkg::PackageDirectory;
14use futures::prelude::*;
15use std::sync::Arc;
16use std::sync::atomic::{AtomicBool, Ordering};
17use zx_status::Status;
18
19#[derive(Debug, Clone)]
21pub struct Client {
22 proxy: fpkg::PackageCacheProxy,
23}
24
25impl Client {
26 pub fn from_proxy(proxy: fpkg::PackageCacheProxy) -> Self {
28 Self { proxy }
29 }
30
31 pub fn proxy(&self) -> &fpkg::PackageCacheProxy {
33 &self.proxy
34 }
35
36 pub fn get(
39 &self,
40 meta_far_blob: BlobInfo,
41 gc_protection: fpkg::GcProtection,
42 ) -> Result<Get, fidl::Error> {
43 let (needed_blobs, needed_blobs_server_end) =
44 fidl::endpoints::create_proxy::<fpkg::NeededBlobsMarker>();
45 let (pkg_dir, pkg_dir_server_end) = PackageDirectory::create_request()?;
46
47 let get_fut = self.proxy.get(
48 &meta_far_blob.into(),
49 gc_protection,
50 needed_blobs_server_end,
51 pkg_dir_server_end,
52 );
53
54 Ok(Get {
55 get_fut,
56 pkg_dir,
57 needed_blobs,
58 pkg_present: SharedBoolEvent::new(),
59 meta_far: meta_far_blob.blob_id,
60 })
61 }
62
63 pub async fn get_already_cached(
73 &self,
74 meta_far_blob: BlobId,
75 ) -> Result<PackageDirectory, GetAlreadyCachedError> {
76 let mut get = self
77 .get(
78 BlobInfo { blob_id: meta_far_blob, length: 0 },
79 fpkg::GcProtection::OpenPackageTracking,
80 )
81 .map_err(GetAlreadyCachedError::Get)?;
82 if let Some(_) = get.open_meta_blob().await.map_err(GetAlreadyCachedError::OpenMetaBlob)? {
83 return Err(GetAlreadyCachedError::MissingMetaFar);
84 }
85
86 if let Some(missing_blobs) = get
87 .get_missing_blobs()
88 .try_next()
89 .await
90 .map_err(GetAlreadyCachedError::GetMissingBlobs)?
91 {
92 return Err(GetAlreadyCachedError::MissingContentBlobs(missing_blobs));
93 }
94
95 get.finish().await.map_err(GetAlreadyCachedError::FinishGet)
96 }
97
98 pub async fn get_subpackage(
101 &self,
102 superpackage: BlobId,
103 subpackage: &fuchsia_url::RelativePackageUrl,
104 ) -> Result<PackageDirectory, GetSubpackageError> {
105 let (dir, dir_server_end) =
106 PackageDirectory::create_request().map_err(GetSubpackageError::CreatingHandles)?;
107 let () = self
108 .proxy
109 .get_subpackage(
110 &superpackage.into(),
111 &fpkg::PackageUrl { url: subpackage.into() },
112 dir_server_end,
113 )
114 .await
115 .map_err(GetSubpackageError::CallingGetSubpackage)??;
116 Ok(dir)
117 }
118
119 pub fn write_blobs(&self) -> Result<WriteBlobs, fidl::Error> {
121 let (needed_blobs, needed_blobs_server_end) =
122 fidl::endpoints::create_proxy::<fpkg::NeededBlobsMarker>();
123
124 let () = self.proxy.write_blobs(needed_blobs_server_end)?;
125
126 Ok(WriteBlobs { needed_blobs })
127 }
128}
129
130#[derive(thiserror::Error, Debug)]
131#[allow(missing_docs)]
132pub enum GetAlreadyCachedError {
133 #[error("calling get")]
134 Get(#[source] fidl::Error),
135
136 #[error("opening meta blob")]
137 OpenMetaBlob(#[source] OpenBlobError),
138
139 #[error("meta.far blob not cached")]
140 MissingMetaFar,
141
142 #[error("getting missing blobs")]
143 GetMissingBlobs(#[source] ListMissingBlobsError),
144
145 #[error("content blobs not cached {0:?}")]
146 MissingContentBlobs(Vec<BlobInfo>),
147
148 #[error("finishing get")]
149 FinishGet(#[source] GetError),
150}
151
152impl GetAlreadyCachedError {
153 pub fn was_not_cached(&self) -> bool {
155 use GetAlreadyCachedError::*;
156 match self {
157 Get(..) | OpenMetaBlob(..) | GetMissingBlobs(..) | FinishGet(..) => false,
158 MissingMetaFar | MissingContentBlobs(..) => true,
159 }
160 }
161}
162
163#[derive(thiserror::Error, Debug)]
164#[allow(missing_docs)]
165pub enum GetSubpackageError {
166 #[error("creating handles")]
167 CreatingHandles(#[source] fidl::Error),
168
169 #[error("calling GetCached FIDL")]
170 CallingGetSubpackage(#[source] fidl::Error),
171
172 #[error("the superpackage does not have an open package connection")]
173 SuperpackageClosed,
174
175 #[error("the subpackage does not exist")]
176 DoesNotExist,
177
178 #[error("internal")]
179 Internal,
180}
181
182impl From<fpkg::GetSubpackageError> for GetSubpackageError {
183 fn from(fidl: fpkg::GetSubpackageError) -> Self {
184 use GetSubpackageError::*;
185 use fpkg::GetSubpackageError as fErr;
186 match fidl {
187 fErr::SuperpackageClosed => SuperpackageClosed,
188 fErr::DoesNotExist => DoesNotExist,
189 fErr::Internal => Internal,
190 }
191 }
192}
193
194#[derive(Debug, Clone)]
195struct SharedBoolEvent(Arc<AtomicBool>);
196
197impl SharedBoolEvent {
198 fn new() -> Self {
199 Self(Arc::new(AtomicBool::new(false)))
200 }
201
202 fn get(&self) -> bool {
203 self.0.load(Ordering::SeqCst)
204 }
205
206 fn set(&self) {
207 self.0.store(true, Ordering::SeqCst)
208 }
209}
210
211async fn open_blob(
212 needed_blobs: &fpkg::NeededBlobsProxy,
213 kind: OpenKind,
214 blob_id: BlobId,
215 pkg_present: Option<&SharedBoolEvent>,
216 allow_existing: bool,
217) -> Result<Option<NeededBlob>, OpenBlobError> {
218 let open_fut = match kind {
219 OpenKind::Meta => needed_blobs.open_meta_blob(),
220 OpenKind::Content => needed_blobs.open_blob(&blob_id.into(), allow_existing),
221 };
222 match open_fut.await {
223 Err(fidl::Error::ClientChannelClosed {
224 epitaph: fidl::Epitaph::Explicit(Ok(())), ..
225 }) => {
226 if let Some(pkg_present) = pkg_present {
227 pkg_present.set();
228 }
229 Ok(None)
230 }
231 res => {
232 if let Some(blob) = res?? {
233 Ok(Some(NeededBlob {
234 blob: Blob {
235 needed_blobs: needed_blobs.clone(),
236 blob_id,
237 state: NeedsTruncate(blob),
238 },
239 }))
240 } else {
241 Ok(None)
242 }
243 }
244 }
245}
246
247#[derive(Debug, Clone, Copy, PartialEq, Eq)]
248enum OpenKind {
249 Meta,
250 Content,
251}
252
253#[derive(Debug)]
255pub struct DeferredOpenBlob {
256 needed_blobs: fpkg::NeededBlobsProxy,
257 kind: OpenKind,
258 blob_id: BlobId,
259 allow_existing: bool,
260 pkg_present: Option<SharedBoolEvent>,
261}
262
263impl DeferredOpenBlob {
264 pub async fn open(&self) -> Result<Option<NeededBlob>, OpenBlobError> {
267 open_blob(
268 &self.needed_blobs,
269 self.kind,
270 self.blob_id,
271 self.pkg_present.as_ref(),
272 self.allow_existing,
273 )
274 .await
275 }
276
277 fn proxy_cmp_key(&self) -> u32 {
278 use fidl::AsHandleRef;
279 use fidl::endpoints::Proxy;
280 self.needed_blobs.as_channel().as_handle_ref().raw_handle()
281 }
282}
283
284impl std::cmp::PartialEq for DeferredOpenBlob {
285 fn eq(&self, other: &Self) -> bool {
286 self.proxy_cmp_key() == other.proxy_cmp_key() && self.kind == other.kind
287 }
288}
289
290impl std::cmp::Eq for DeferredOpenBlob {}
291
292#[derive(Debug)]
298pub struct Get {
299 get_fut: fidl::client::QueryResponseFut<Result<(), i32>>,
300 needed_blobs: fpkg::NeededBlobsProxy,
301 pkg_dir: PackageDirectory,
302 pkg_present: SharedBoolEvent,
303 meta_far: BlobId,
304}
305
306impl Get {
307 pub fn make_open_meta_blob(&mut self) -> DeferredOpenBlob {
310 DeferredOpenBlob {
311 needed_blobs: self.needed_blobs.clone(),
312 kind: OpenKind::Meta,
313 blob_id: self.meta_far,
314 allow_existing: false,
315 pkg_present: Some(self.pkg_present.clone()),
316 }
317 }
318
319 pub async fn open_meta_blob(&mut self) -> Result<Option<NeededBlob>, OpenBlobError> {
322 open_blob(&self.needed_blobs, OpenKind::Meta, self.meta_far, Some(&self.pkg_present), false)
323 .await
324 }
325
326 fn start_get_missing_blobs(
327 &mut self,
328 ) -> Result<Option<fpkg::BlobInfoIteratorProxy>, fidl::Error> {
329 if self.pkg_present.get() {
330 return Ok(None);
331 }
332
333 let (blob_iterator, blob_iterator_server_end) =
334 fidl::endpoints::create_proxy::<fpkg::BlobInfoIteratorMarker>();
335
336 self.needed_blobs.get_missing_blobs(blob_iterator_server_end)?;
337 Ok(Some(blob_iterator))
338 }
339
340 pub fn get_missing_blobs(
346 &mut self,
347 ) -> impl Stream<Item = Result<Vec<BlobInfo>, ListMissingBlobsError>> + Unpin + use<> {
348 match self.start_get_missing_blobs() {
349 Ok(option_iter) => match option_iter {
350 Some(iterator) => crate::fidl_iterator_to_stream(iterator)
351 .map_ok(|v| v.into_iter().map(BlobInfo::from).collect())
352 .map_err(ListMissingBlobsError::CallNextOnBlobIterator)
353 .left_stream(),
354 None => futures::stream::empty().right_stream(),
355 }
356 .left_stream(),
357 Err(e) => {
358 futures::stream::iter(Some(Err(ListMissingBlobsError::CallGetMissingBlobs(e))))
359 .right_stream()
360 }
361 }
362 }
363
364 pub fn make_open_blob(&mut self, content_blob: BlobId) -> DeferredOpenBlob {
367 DeferredOpenBlob {
368 needed_blobs: self.needed_blobs.clone(),
369 kind: OpenKind::Content,
370 blob_id: content_blob,
371 allow_existing: false,
372 pkg_present: None,
373 }
374 }
375
376 pub async fn open_blob(
379 &mut self,
380 content_blob: BlobId,
381 ) -> Result<Option<NeededBlob>, OpenBlobError> {
382 open_blob(&self.needed_blobs, OpenKind::Content, content_blob, None, false).await
383 }
384
385 pub async fn finish(self) -> Result<PackageDirectory, GetError> {
388 drop(self.needed_blobs);
389 let () = self.get_fut.await?.map_err(Status::from_raw)?;
390 Ok(self.pkg_dir)
391 }
392
393 pub async fn abort(self) {
395 self.needed_blobs.abort().map(|_: Result<(), fidl::Error>| ()).await;
396 let _ = self.get_fut.await;
401 }
402}
403
404#[derive(Clone, Debug)]
406pub struct WriteBlobs {
407 needed_blobs: fpkg::NeededBlobsProxy,
408}
409
410impl WriteBlobs {
411 pub fn make_open_blob(&mut self, blob: BlobId, allow_existing: bool) -> DeferredOpenBlob {
414 DeferredOpenBlob {
415 needed_blobs: self.needed_blobs.clone(),
416 kind: OpenKind::Content,
417 blob_id: blob,
418 allow_existing,
419 pkg_present: None,
420 }
421 }
422
423 pub async fn open_blob(
425 &mut self,
426 blob: BlobId,
427 allow_existing: bool,
428 ) -> Result<Option<NeededBlob>, OpenBlobError> {
429 open_blob(&self.needed_blobs, OpenKind::Content, blob, None, allow_existing).await
430 }
431}
432
433#[derive(Debug)]
435pub struct NeededBlob {
436 pub blob: Blob<NeedsTruncate>,
439}
440
441#[derive(Debug)]
443pub enum TruncateBlobSuccess {
444 NeedsData(Blob<NeedsData>),
446
447 AllWritten(Blob<NeedsBlobWritten>),
450}
451
452#[derive(Debug)]
454pub enum BlobWriteSuccess {
455 NeedsData(Blob<NeedsData>),
457
458 AllWritten(Blob<NeedsBlobWritten>),
461}
462
463#[derive(Debug)]
465pub struct NeedsTruncate(fidl::endpoints::ClientEnd<ffxfs::BlobWriterMarker>);
466
467#[derive(Debug)]
469pub struct NeedsData {
470 size: u64,
471 written: u64,
472 writer: blob_writer::BlobWriter,
473}
474
475#[derive(Debug)]
478pub struct NeedsBlobWritten;
479
480#[derive(Debug)]
482#[must_use]
483pub struct Blob<S> {
484 needed_blobs: fpkg::NeededBlobsProxy,
485 blob_id: BlobId,
486 state: S,
487}
488
489impl Blob<NeedsTruncate> {
490 pub async fn truncate(self, size: u64) -> Result<TruncateBlobSuccess, TruncateBlobError> {
492 let Self { needed_blobs, blob_id, state: NeedsTruncate(blob) } = self;
493
494 let writer = blob_writer::BlobWriter::create(blob.into_proxy(), size).await.map_err(
495 |e| match e {
496 blob_writer::CreateError::GetVmo(s) if s == Status::NO_SPACE => {
497 TruncateBlobError::NoSpace
498 }
499 _ => TruncateBlobError::CreateBlobWriter(e),
500 },
501 )?;
502
503 Ok(if size == 0 {
504 TruncateBlobSuccess::AllWritten(Blob { needed_blobs, blob_id, state: NeedsBlobWritten })
505 } else {
506 TruncateBlobSuccess::NeedsData(Blob {
507 needed_blobs,
508 blob_id,
509 state: NeedsData { size, written: 0, writer },
510 })
511 })
512 }
513
514 pub fn deconstruct(
522 self,
523 ) -> (fpkg::NeededBlobsProxy, BlobId, fidl::endpoints::ClientEnd<ffxfs::BlobWriterMarker>) {
524 let Blob { needed_blobs, blob_id, state: NeedsTruncate(client) } = self;
525 (needed_blobs, blob_id, client)
526 }
527}
528
529impl Blob<NeedsData> {
530 pub fn write(
536 self,
537 buf: &[u8],
538 ) -> impl Future<Output = Result<BlobWriteSuccess, WriteBlobError>> + '_ {
539 self.write_with_trace_callbacks(buf, &|_| {}, &|| {})
540 }
541
542 pub async fn write_with_trace_callbacks(
554 mut self,
555 buf: &[u8],
556 after_write: &(dyn Fn(u64) + Send + Sync),
557 after_write_ack: &(dyn Fn() + Send + Sync),
558 ) -> Result<BlobWriteSuccess, WriteBlobError> {
559 assert!(self.state.written + buf.len() as u64 <= self.state.size);
560
561 let fut = self.state.writer.write(buf);
562 let () = after_write(buf.len() as u64);
563 let res = fut.await;
564 let () = after_write_ack();
565 let () = res.map_err(|e| match e {
566 e @ blob_writer::WriteError::BytesReady(s) => match s {
567 Status::IO_DATA_INTEGRITY => WriteBlobError::Corrupt,
568 Status::NO_SPACE => WriteBlobError::NoSpace,
569 _ => WriteBlobError::FxBlob(e),
570 },
571 e => WriteBlobError::FxBlob(e),
572 })?;
573
574 self.state.written += buf.len() as u64;
575
576 if self.state.written == self.state.size {
577 let Self { needed_blobs, blob_id, state: _ } = self;
578 Ok(BlobWriteSuccess::AllWritten(Blob {
579 needed_blobs,
580 blob_id,
581 state: NeedsBlobWritten,
582 }))
583 } else {
584 Ok(BlobWriteSuccess::NeedsData(self))
585 }
586 }
587}
588
589impl Blob<NeedsBlobWritten> {
590 pub async fn blob_written(self) -> Result<(), BlobWrittenError> {
592 Ok(self.needed_blobs.blob_written(&self.blob_id.into()).await??)
593 }
594}
595
596#[derive(Debug, thiserror::Error)]
598#[allow(missing_docs)]
599pub enum OpenError {
600 #[error("the package does not exist")]
601 NotFound,
602
603 #[error("Open() responded with an unexpected status")]
604 UnexpectedResponse(#[source] Status),
605
606 #[error("transport error")]
607 Fidl(#[from] fidl::Error),
608}
609#[derive(Debug, thiserror::Error)]
611#[allow(missing_docs)]
612pub enum GetError {
613 #[error("Get() responded with an unexpected status")]
614 UnexpectedResponse(#[from] Status),
615
616 #[error("transport error")]
617 Fidl(#[from] fidl::Error),
618}
619
620#[derive(Debug, thiserror::Error)]
622#[allow(missing_docs)]
623pub enum OpenBlobError {
624 #[error("there is insufficient storage space available to persist this blob")]
625 OutOfSpace,
626
627 #[error("this blob is already open for write by another cache operation")]
628 ConcurrentWrite,
629
630 #[error("an unspecified error occurred during underlying I/O")]
631 UnspecifiedIo,
632
633 #[error("an unspecified error occurred")]
634 Internal,
635
636 #[error("transport error")]
637 Fidl(#[from] fidl::Error),
638}
639
640impl From<fpkg::OpenBlobError> for OpenBlobError {
641 fn from(e: fpkg::OpenBlobError) -> Self {
642 match e {
643 fpkg::OpenBlobError::OutOfSpace => OpenBlobError::OutOfSpace,
644 fpkg::OpenBlobError::ConcurrentWrite => OpenBlobError::ConcurrentWrite,
645 fpkg::OpenBlobError::UnspecifiedIo => OpenBlobError::UnspecifiedIo,
646 fpkg::OpenBlobError::Internal => OpenBlobError::Internal,
647 }
648 }
649}
650
651#[derive(Debug, thiserror::Error)]
653#[allow(missing_docs)]
654pub enum ListMissingBlobsError {
655 #[error("while obtaining the missing blobs fidl iterator")]
656 CallGetMissingBlobs(#[source] fidl::Error),
657
658 #[error("while obtaining the next chunk of blobs from the fidl iterator")]
659 CallNextOnBlobIterator(#[source] fidl::Error),
660}
661
662#[derive(Debug, thiserror::Error)]
664#[allow(missing_docs)]
665pub enum TruncateBlobError {
666 #[error("insufficient storage space is available")]
667 NoSpace,
668
669 #[error("creating blob writer")]
670 CreateBlobWriter(#[source] blob_writer::CreateError),
671
672 #[error("transport error")]
673 Fidl(#[from] fidl::Error),
674
675 #[error("blob is in an invalid state")]
676 BadState,
677}
678
679#[derive(Debug, thiserror::Error)]
681#[allow(missing_docs)]
682pub enum WriteBlobError {
683 #[error("the written data was corrupt")]
684 Corrupt,
685
686 #[error("insufficient storage space is available")]
687 NoSpace,
688
689 #[error("transport error")]
690 Fidl(#[from] fidl::Error),
691
692 #[error("while using the fxblob writer")]
693 FxBlob(#[source] blob_writer::WriteError),
694}
695
696#[derive(Debug, thiserror::Error)]
698#[allow(missing_docs)]
699pub enum BlobWrittenError {
700 #[error("pkg-cache could not find the blob after it was successfully written")]
701 MissingAfterWritten,
702
703 #[error("NeededBlobs.BlobWritten was called before the blob was opened")]
704 UnopenedBlob,
705
706 #[error("transport error")]
707 Fidl(#[from] fidl::Error),
708}
709
710impl From<fpkg::BlobWrittenError> for BlobWrittenError {
711 fn from(e: fpkg::BlobWrittenError) -> Self {
712 match e {
713 fpkg::BlobWrittenError::NotWritten => BlobWrittenError::MissingAfterWritten,
714 fpkg::BlobWrittenError::UnopenedBlob => BlobWrittenError::UnopenedBlob,
715 }
716 }
717}
718
719#[cfg(test)]
720mod tests {
721 use super::*;
722 use assert_matches::assert_matches;
723 use fidl::endpoints::{ClientEnd, RequestStream as _};
724 use fidl_fuchsia_io as fio;
725 use fidl_fuchsia_pkg::{
726 BlobInfoIteratorRequest, NeededBlobsRequest, NeededBlobsRequestStream,
727 PackageCacheGetResponder, PackageCacheMarker, PackageCacheRequest,
728 PackageCacheRequestStream,
729 };
730
731 struct MockPackageCache {
732 stream: PackageCacheRequestStream,
733 }
734
735 impl MockPackageCache {
736 fn new() -> (Client, Self) {
737 let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<PackageCacheMarker>();
738 (Client::from_proxy(proxy), Self { stream })
739 }
740
741 async fn expect_get(
742 &mut self,
743 blob_info: BlobInfo,
744 expected_gc_protection: fpkg::GcProtection,
745 ) -> PendingGet {
746 match self.stream.next().await {
747 Some(Ok(PackageCacheRequest::Get {
748 meta_far_blob,
749 gc_protection,
750 needed_blobs,
751 dir,
752 responder,
753 })) => {
754 assert_eq!(BlobInfo::from(meta_far_blob), blob_info);
755 assert_eq!(gc_protection, expected_gc_protection);
756 let needed_blobs = needed_blobs.into_stream();
757 let dir = dir.into_stream();
758
759 PendingGet { stream: needed_blobs, dir, responder }
760 }
761 r => panic!("Unexpected request: {r:?}"),
762 }
763 }
764
765 async fn expect_closed(mut self) {
766 assert_matches!(self.stream.next().await, None);
767 }
768 }
769
770 struct PendingGet {
771 stream: NeededBlobsRequestStream,
772 dir: fio::DirectoryRequestStream,
773 responder: PackageCacheGetResponder,
774 }
775
776 impl PendingGet {
777 async fn new() -> (Get, PendingGet) {
778 let (client, mut server) = MockPackageCache::new();
779
780 let get = client.get(blob_info(42), fpkg::GcProtection::OpenPackageTracking).unwrap();
781 let pending_get =
782 server.expect_get(blob_info(42), fpkg::GcProtection::OpenPackageTracking).await;
783 (get, pending_get)
784 }
785
786 fn finish_hold_stream_open(self) -> (NeededBlobsRequestStream, PackageDirProvider) {
787 self.stream.control_handle().shutdown_with_epitaph(Ok(()));
788 self.responder.send(Ok(())).unwrap();
789 (self.stream, PackageDirProvider { stream: self.dir })
790 }
791
792 fn finish(self) -> PackageDirProvider {
793 self.stream.control_handle().shutdown_with_epitaph(Ok(()));
794 self.responder.send(Ok(())).unwrap();
795 PackageDirProvider { stream: self.dir }
796 }
797
798 #[cfg(target_os = "fuchsia")]
799 fn fail_the_get(self) {
800 self.responder
801 .send(Err(Status::IO_INVALID.into_raw()))
802 .expect("client should be waiting");
803 }
804
805 async fn expect_open_meta_blob(
806 mut self,
807 res: Result<Option<ClientEnd<ffxfs::BlobWriterMarker>>, fpkg::OpenBlobError>,
808 ) -> Self {
809 match self.stream.next().await {
810 Some(Ok(NeededBlobsRequest::OpenMetaBlob { responder })) => {
811 responder.send(res).unwrap();
812 }
813 r => panic!("Unexpected request: {r:?}"),
814 }
815 self
816 }
817
818 async fn expect_open_blob(
819 mut self,
820 expected_blob_id: BlobId,
821 res: Result<Option<ClientEnd<ffxfs::BlobWriterMarker>>, fpkg::OpenBlobError>,
822 ) -> Self {
823 match self.stream.next().await {
824 Some(Ok(NeededBlobsRequest::OpenBlob {
825 blob_id,
826 allow_existing: _,
827 responder,
828 })) => {
829 assert_eq!(BlobId::from(blob_id), expected_blob_id);
830 responder.send(res).unwrap();
831 }
832 r => panic!("Unexpected request: {r:?}"),
833 }
834 self
835 }
836
837 async fn expect_get_missing_blobs(mut self, response_chunks: Vec<Vec<BlobInfo>>) -> Self {
838 match self.stream.next().await {
839 Some(Ok(NeededBlobsRequest::GetMissingBlobs { iterator, control_handle: _ })) => {
840 let mut stream = iterator.into_stream();
841
842 for chunk in response_chunks {
844 let chunk = chunk.into_iter().map(fpkg::BlobInfo::from).collect::<Vec<_>>();
845
846 let BlobInfoIteratorRequest::Next { responder } =
847 stream.next().await.unwrap().unwrap();
848 responder.send(&chunk).unwrap();
849 }
850
851 let BlobInfoIteratorRequest::Next { responder } =
853 stream.next().await.unwrap().unwrap();
854 responder.send(&[]).unwrap();
855
856 assert_matches!(stream.next().await, None);
858 }
859 r => panic!("Unexpected request: {r:?}"),
860 }
861 self
862 }
863
864 async fn expect_get_missing_blobs_client_closes_channel(
865 mut self,
866 response_chunks: Vec<Vec<BlobInfo>>,
867 ) -> Self {
868 match self.stream.next().await {
869 Some(Ok(NeededBlobsRequest::GetMissingBlobs { iterator, control_handle: _ })) => {
870 let mut stream = iterator.into_stream();
871
872 for chunk in response_chunks {
874 let chunk = chunk.into_iter().map(fpkg::BlobInfo::from).collect::<Vec<_>>();
875
876 let BlobInfoIteratorRequest::Next { responder } =
877 stream.next().await.unwrap().unwrap();
878 responder.send(&chunk).unwrap();
879 }
880
881 assert_matches!(stream.next().await, None);
883 }
884 r => panic!("Unexpected request: {r:?}"),
885 }
886 self
887 }
888
889 async fn expect_get_missing_blobs_inject_iterator_error(mut self) -> Self {
890 match self.stream.next().await {
891 Some(Ok(NeededBlobsRequest::GetMissingBlobs { iterator, control_handle: _ })) => {
892 iterator
893 .into_stream_and_control_handle()
894 .1
895 .shutdown_with_epitaph(Status::ADDRESS_IN_USE);
896 }
897 r => panic!("Unexpected request: {r:?}"),
898 }
899 self
900 }
901
902 #[cfg(target_os = "fuchsia")]
903 async fn expect_abort(mut self) -> Self {
904 match self.stream.next().await {
905 Some(Ok(NeededBlobsRequest::Abort { responder })) => {
906 responder.send().unwrap();
907 }
908 r => panic!("Unexpected request: {r:?}"),
909 }
910 self
911 }
912 }
913
914 struct PackageDirProvider {
915 stream: fio::DirectoryRequestStream,
916 }
917
918 impl PackageDirProvider {
919 fn close_pkg_dir(self) {
920 self.stream.control_handle().shutdown_with_epitaph(Status::NOT_EMPTY);
921 }
922 }
923
924 fn blob_id(n: u8) -> BlobId {
925 BlobId::from([n; 32])
926 }
927
928 fn blob_info(n: u8) -> BlobInfo {
929 BlobInfo { blob_id: blob_id(n), length: 0 }
930 }
931
932 #[fuchsia::test]
933 async fn constructor() {
934 let (proxy, stream) = fidl::endpoints::create_proxy_and_stream::<PackageCacheMarker>();
935 let client = Client::from_proxy(proxy);
936
937 drop(stream);
938 assert_matches!(client.proxy().sync().await, Err(_));
939 }
940
941 #[fuchsia::test]
942 async fn get_present_package() {
943 let (client, mut server) = MockPackageCache::new();
944
945 let ((), ()) = future::join(
946 async {
947 server
948 .expect_get(blob_info(2), fpkg::GcProtection::OpenPackageTracking)
949 .await
950 .finish()
951 .close_pkg_dir();
952 server.expect_closed().await;
953 },
954 async move {
955 let mut get =
956 client.get(blob_info(2), fpkg::GcProtection::OpenPackageTracking).unwrap();
957
958 assert_matches!(get.open_meta_blob().await.unwrap(), None);
959 assert_eq!(get.get_missing_blobs().try_concat().await.unwrap(), vec![]);
960 let pkg_dir = get.finish().await.unwrap();
961
962 assert_matches!(
963 pkg_dir.into_proxy().take_event_stream().next().await,
964 Some(Err(fidl::Error::ClientChannelClosed { epitaph, .. }))
965 if epitaph == Status::NOT_EMPTY
966 );
967 },
968 )
969 .await;
970 }
971
972 #[fuchsia::test]
973 async fn get_present_package_handles_slow_stream_close() {
974 let (client, mut server) = MockPackageCache::new();
975
976 let (send, recv) = futures::channel::oneshot::channel::<()>();
977
978 let ((), ()) = future::join(
979 async {
980 let (needed_blobs_stream, pkg_dir) = server
981 .expect_get(blob_info(2), fpkg::GcProtection::OpenPackageTracking)
982 .await
983 .finish_hold_stream_open();
984 pkg_dir.close_pkg_dir();
985
986 let _ = recv.await;
988 drop(needed_blobs_stream);
989 },
990 async move {
991 let mut get =
992 client.get(blob_info(2), fpkg::GcProtection::OpenPackageTracking).unwrap();
993
994 assert_matches!(get.open_meta_blob().await.unwrap(), None);
995
996 let missing_blobs_stream = get.get_missing_blobs();
1000 drop(send);
1001 assert_eq!(missing_blobs_stream.try_concat().await.unwrap(), vec![]);
1002 let pkg_dir = get.finish().await.unwrap();
1003
1004 assert_matches!(
1005 pkg_dir.into_proxy().take_event_stream().next().await,
1006 Some(Err(fidl::Error::ClientChannelClosed { epitaph, .. }))
1007 if epitaph == Status::NOT_EMPTY
1008 );
1009 },
1010 )
1011 .await;
1012 }
1013
1014 #[fuchsia::test]
1015 async fn needed_blobs_open_meta_far() {
1016 let (mut get, pending_get) = PendingGet::new().await;
1017
1018 let ((), ()) = future::join(
1019 async {
1020 pending_get
1021 .expect_open_meta_blob(Ok(None))
1022 .await
1023 .expect_open_meta_blob(Ok(Some(fidl::endpoints::create_endpoints().0)))
1024 .await
1025 .expect_open_meta_blob(Ok(None))
1026 .await
1027 .expect_open_meta_blob(Ok(Some(fidl::endpoints::create_endpoints().0)))
1028 .await
1029 .expect_open_meta_blob(Err(fpkg::OpenBlobError::OutOfSpace))
1030 .await
1031 .expect_open_meta_blob(Err(fpkg::OpenBlobError::UnspecifiedIo))
1032 .await;
1033 },
1034 async {
1035 {
1036 let opener = get.make_open_meta_blob();
1037 assert_matches!(opener.open().await.unwrap(), None);
1038 assert_matches!(opener.open().await.unwrap(), Some(_));
1039 }
1040 assert_matches!(get.open_meta_blob().await.unwrap(), None);
1041 assert_matches!(get.open_meta_blob().await.unwrap(), Some(_));
1042 assert_matches!(get.open_meta_blob().await, Err(OpenBlobError::OutOfSpace));
1043 assert_matches!(get.open_meta_blob().await, Err(OpenBlobError::UnspecifiedIo));
1044 },
1045 )
1046 .await;
1047 }
1048
1049 #[fuchsia::test]
1050 async fn needed_blobs_open_content_blob() {
1051 let (mut get, pending_get) = PendingGet::new().await;
1052
1053 let ((), ()) = future::join(
1054 async {
1055 pending_get
1056 .expect_open_blob(blob_id(2), Ok(None))
1057 .await
1058 .expect_open_blob(blob_id(2), Ok(Some(fidl::endpoints::create_endpoints().0)))
1059 .await
1060 .expect_open_blob(blob_id(10), Ok(None))
1061 .await
1062 .expect_open_blob(blob_id(11), Ok(Some(fidl::endpoints::create_endpoints().0)))
1063 .await
1064 .expect_open_blob(blob_id(12), Err(fpkg::OpenBlobError::OutOfSpace))
1065 .await
1066 .expect_open_blob(blob_id(13), Err(fpkg::OpenBlobError::UnspecifiedIo))
1067 .await;
1068 },
1069 async {
1070 {
1071 let opener = get.make_open_blob(blob_id(2));
1072 assert_matches!(opener.open().await.unwrap(), None);
1073 assert_matches!(opener.open().await.unwrap(), Some(_));
1074 }
1075 assert_matches!(get.open_blob(blob_id(10),).await.unwrap(), None);
1076 assert_matches!(get.open_blob(blob_id(11),).await.unwrap(), Some(_));
1077 assert_matches!(get.open_blob(blob_id(12),).await, Err(OpenBlobError::OutOfSpace));
1078 assert_matches!(
1079 get.open_blob(blob_id(13),).await,
1080 Err(OpenBlobError::UnspecifiedIo)
1081 );
1082 },
1083 )
1084 .await;
1085 }
1086
1087 #[fuchsia::test]
1088 async fn needed_blobs_get_missing_blobs_on_closed_ok() {
1089 let (mut get, pending_get) = PendingGet::new().await;
1090 let _ = pending_get.finish();
1091
1092 assert_matches!(get.open_meta_blob().await, Ok(None));
1093 assert_eq!(get.get_missing_blobs().try_concat().await.unwrap(), vec![]);
1094 }
1095
1096 #[fuchsia::test]
1097 async fn needed_blobs_get_missing_blobs() {
1098 let (mut get, pending_get) = PendingGet::new().await;
1099
1100 let ((), ()) = future::join(
1101 async {
1102 pending_get
1103 .expect_get_missing_blobs(vec![
1104 vec![blob_info(1), blob_info(2)],
1105 vec![blob_info(3)],
1106 ])
1107 .await;
1108 },
1109 async {
1110 assert_eq!(
1111 get.get_missing_blobs().try_concat().await.unwrap(),
1112 vec![blob_info(1), blob_info(2), blob_info(3)]
1113 );
1114 },
1115 )
1116 .await;
1117 }
1118
1119 #[fuchsia::test]
1120 async fn needed_blobs_get_missing_blobs_fail_to_obtain_iterator() {
1121 let (mut get, pending_get) = PendingGet::new().await;
1122 drop(pending_get);
1123
1124 assert_matches!(
1125 get.get_missing_blobs().try_concat().await,
1126 Err(ListMissingBlobsError::CallNextOnBlobIterator(
1127 fidl::Error::ClientChannelClosed { epitaph, ..})
1128 )
1129 if epitaph == Status::PEER_CLOSED
1130 );
1131 }
1132
1133 #[fuchsia::test]
1134 async fn needed_blobs_get_missing_blobs_iterator_contains_error() {
1135 let (mut get, pending_get) = PendingGet::new().await;
1136
1137 let (_, ()) =
1138 future::join(pending_get.expect_get_missing_blobs_inject_iterator_error(), async {
1139 assert_matches!(
1140 get.get_missing_blobs().try_concat().await,
1141 Err(ListMissingBlobsError::CallNextOnBlobIterator(
1142 fidl::Error::ClientChannelClosed { epitaph, ..}
1143 ))
1144 if epitaph == Status::ADDRESS_IN_USE
1145 );
1146 })
1147 .await;
1148 }
1149
1150 #[cfg(target_os = "fuchsia")]
1151 #[test]
1152 fn needed_blobs_abort() {
1153 use futures::future::Either;
1154 use futures::pin_mut;
1155 use std::task::Poll;
1156
1157 let mut executor = fuchsia_async::TestExecutor::new_with_fake_time();
1158
1159 let fut = async {
1160 let (get, pending_get) = PendingGet::new().await;
1161
1162 let abort_fut = get.abort().boxed();
1163 let expect_abort_fut = pending_get.expect_abort();
1164 pin_mut!(expect_abort_fut);
1165
1166 match futures::future::select(abort_fut, expect_abort_fut).await {
1167 Either::Left(((), _expect_abort_fut)) => {
1168 panic!("abort should wait for the get future to complete")
1169 }
1170 Either::Right((pending_get, abort_fut)) => (abort_fut, pending_get),
1171 }
1172 };
1173 pin_mut!(fut);
1174
1175 let (mut abort_fut, pending_get) = match executor.run_until_stalled(&mut fut) {
1176 Poll::Pending => panic!("should complete"),
1177 Poll::Ready((abort_fut, pending_get)) => (abort_fut, pending_get),
1178 };
1179
1180 assert_matches!(executor.run_until_stalled(&mut abort_fut), Poll::Pending);
1182 pending_get.fail_the_get();
1183 assert_matches!(executor.run_until_stalled(&mut abort_fut), Poll::Ready(()));
1184 }
1185
1186 struct MockNeededBlob {
1187 writer: ffxfs::BlobWriterRequestStream,
1188 needed_blobs: fpkg::NeededBlobsRequestStream,
1189 vmo: Option<zx::Vmo>,
1190 }
1191
1192 impl MockNeededBlob {
1193 fn mock_hash() -> BlobId {
1194 [7; 32].into()
1195 }
1196
1197 fn new() -> (NeededBlob, Self) {
1198 let (writer_client, writer) =
1199 fidl::endpoints::create_request_stream::<ffxfs::BlobWriterMarker>();
1200 let (needed_blobs_proxy, needed_blobs) =
1201 fidl::endpoints::create_proxy_and_stream::<fpkg::NeededBlobsMarker>();
1202 (
1203 NeededBlob {
1204 blob: Blob {
1205 needed_blobs: needed_blobs_proxy,
1206 blob_id: Self::mock_hash(),
1207 state: NeedsTruncate(writer_client),
1208 },
1209 },
1210 Self { writer, needed_blobs, vmo: None },
1211 )
1212 }
1213
1214 async fn fail_get_vmo(mut self) -> Self {
1215 match self.writer.next().await {
1216 Some(Ok(ffxfs::BlobWriterRequest::GetVmo { size: _, responder })) => {
1217 responder.send(Err(Status::NO_SPACE.into_raw())).unwrap();
1218 }
1219 r => panic!("Unexpected request: {r:?}"),
1220 }
1221 self
1222 }
1223
1224 async fn expect_get_vmo(mut self, expected_size: u64) -> Self {
1225 match self.writer.next().await {
1226 Some(Ok(ffxfs::BlobWriterRequest::GetVmo { size, responder })) => {
1227 assert_eq!(size, expected_size);
1228 let vmo = zx::Vmo::create(size).unwrap();
1229 assert!(self.vmo.is_none());
1230 self.vmo = Some(vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap());
1231 responder.send(Ok(vmo)).unwrap();
1232 }
1233 r => panic!("Unexpected request: {r:?}"),
1234 }
1235 self
1236 }
1237
1238 async fn fail_bytes_ready(mut self) -> Self {
1239 match self.writer.next().await {
1240 Some(Ok(ffxfs::BlobWriterRequest::BytesReady { bytes_written: _, responder })) => {
1241 responder.send(Err(Status::NO_SPACE.into_raw())).unwrap();
1242 }
1243 r => panic!("Unexpected request: {r:?}"),
1244 }
1245 self
1246 }
1247
1248 async fn expect_bytes_ready(mut self, expected_payload: &[u8], offset: u64) -> Self {
1249 match self.writer.next().await {
1250 Some(Ok(ffxfs::BlobWriterRequest::BytesReady { bytes_written, responder })) => {
1251 assert_eq!(bytes_written, u64::try_from(expected_payload.len()).unwrap());
1252
1253 let vmo = self.vmo.as_ref().unwrap();
1254 let mut buf = vec![0; expected_payload.len()];
1255 let () = vmo.read(&mut buf, offset).unwrap();
1256 assert_eq!(buf, expected_payload);
1257
1258 responder.send(Ok(())).unwrap();
1259 }
1260 r => panic!("Unexpected request: {r:?}"),
1261 }
1262 self
1263 }
1264
1265 async fn expect_blob_written(mut self) -> Self {
1266 match self.needed_blobs.next().await {
1267 Some(Ok(fpkg::NeededBlobsRequest::BlobWritten { blob_id, responder })) => {
1268 assert_eq!(blob_id, Self::mock_hash().into());
1269 responder.send(Ok(())).unwrap();
1270 }
1271 r => panic!("Unexpected request: {r:?}"),
1272 }
1273 self
1274 }
1275 }
1276
1277 #[fuchsia::test]
1278 async fn empty_blob_write() {
1279 let (NeededBlob { blob }, blob_server) = MockNeededBlob::new();
1280
1281 let ((), ()) = future::join(
1282 async {
1283 blob_server.expect_get_vmo(0).await.expect_blob_written().await;
1284 },
1285 async {
1286 let blob = match blob.truncate(0).await.unwrap() {
1287 TruncateBlobSuccess::AllWritten(blob) => blob,
1288 other => panic!("empty blob shouldn't need bytes {other:?}"),
1289 };
1290 let () = blob.blob_written().await.unwrap();
1291 },
1292 )
1293 .await;
1294 }
1295
1296 impl TruncateBlobSuccess {
1297 fn unwrap_needs_data(self) -> Blob<NeedsData> {
1298 match self {
1299 TruncateBlobSuccess::NeedsData(blob) => blob,
1300 TruncateBlobSuccess::AllWritten(_) => panic!("blob should need data"),
1301 }
1302 }
1303 }
1304
1305 impl BlobWriteSuccess {
1306 fn unwrap_needs_data(self) -> Blob<NeedsData> {
1307 match self {
1308 BlobWriteSuccess::NeedsData(blob) => blob,
1309 BlobWriteSuccess::AllWritten(_) => panic!("blob should need data"),
1310 }
1311 }
1312
1313 fn unwrap_all_written(self) -> Blob<NeedsBlobWritten> {
1314 match self {
1315 BlobWriteSuccess::NeedsData(_) => panic!("blob should be completely written"),
1316 BlobWriteSuccess::AllWritten(blob) => blob,
1317 }
1318 }
1319 }
1320
1321 #[fuchsia::test]
1322 async fn small_blob_write() {
1323 let (NeededBlob { blob }, blob_server) = MockNeededBlob::new();
1324
1325 let ((), ()) = future::join(
1326 async {
1327 blob_server
1328 .expect_get_vmo(4)
1329 .await
1330 .expect_bytes_ready(b"test", 0)
1331 .await
1332 .expect_blob_written()
1333 .await;
1334 },
1335 async {
1336 let blob = blob.truncate(4).await.unwrap().unwrap_needs_data();
1337 let blob = blob.write(b"test").await.unwrap().unwrap_all_written();
1338 let () = blob.blob_written().await.unwrap();
1339 },
1340 )
1341 .await;
1342 }
1343
1344 #[fuchsia::test]
1345 async fn blob_truncate_no_space() {
1346 let (NeededBlob { blob }, blob_server) = MockNeededBlob::new();
1347
1348 let ((), ()) = future::join(
1349 async {
1350 blob_server.fail_get_vmo().await;
1351 },
1352 async {
1353 assert_matches!(blob.truncate(4).await, Err(TruncateBlobError::NoSpace));
1354 },
1355 )
1356 .await;
1357 }
1358
1359 #[fuchsia::test]
1360 async fn blob_write_no_space() {
1361 let (NeededBlob { blob }, blob_server) = MockNeededBlob::new();
1362
1363 let ((), ()) = future::join(
1364 async {
1365 blob_server.expect_get_vmo(4).await.fail_bytes_ready().await;
1366 },
1367 async {
1368 let blob = blob.truncate(4).await.unwrap().unwrap_needs_data();
1369 assert_matches!(blob.write(b"test").await, Err(WriteBlobError::NoSpace));
1370 },
1371 )
1372 .await;
1373 }
1374
1375 #[fuchsia::test]
1376 async fn blob_write_multiple_write() {
1377 let (NeededBlob { blob }, blob_server) = MockNeededBlob::new();
1378
1379 let ((), ()) = future::join(
1380 async {
1381 blob_server
1382 .expect_get_vmo(6)
1383 .await
1384 .expect_bytes_ready(b"abc", 0)
1385 .await
1386 .expect_bytes_ready(b"123", 3)
1387 .await
1388 .expect_blob_written()
1389 .await;
1390 },
1391 async {
1392 let blob = blob.truncate(6).await.unwrap().unwrap_needs_data();
1393 let blob = blob.write(b"abc").await.unwrap().unwrap_needs_data();
1394 let blob = blob.write(b"123").await.unwrap().unwrap_all_written();
1395 let () = blob.blob_written().await.unwrap();
1396 },
1397 )
1398 .await;
1399 }
1400
1401 #[fuchsia::test]
1402 async fn get_already_cached_success() {
1403 let (client, mut server) = MockPackageCache::new();
1404
1405 let ((), ()) = future::join(
1406 async {
1407 server
1408 .expect_get(blob_info(2), fpkg::GcProtection::OpenPackageTracking)
1409 .await
1410 .finish()
1411 .close_pkg_dir();
1412 server.expect_closed().await;
1413 },
1414 async move {
1415 let pkg_dir = client.get_already_cached(blob_id(2)).await.unwrap();
1416
1417 assert_matches!(
1418 pkg_dir.into_proxy().take_event_stream().next().await,
1419 Some(Err(fidl::Error::ClientChannelClosed { epitaph, .. }))
1420 if epitaph == Status::NOT_EMPTY
1421 );
1422 },
1423 )
1424 .await;
1425 }
1426
1427 #[fuchsia::test]
1428 async fn get_already_cached_missing_meta_far() {
1429 let (client, mut server) = MockPackageCache::new();
1430
1431 let ((), ()) = future::join(
1432 async {
1433 server
1434 .expect_get(blob_info(2), fpkg::GcProtection::OpenPackageTracking)
1435 .await
1436 .expect_open_meta_blob(Ok(Some(fidl::endpoints::create_endpoints().0)))
1437 .await;
1438 },
1439 async move {
1440 assert_matches!(
1441 client.get_already_cached(blob_id(2)).await,
1442 Err(GetAlreadyCachedError::MissingMetaFar)
1443 );
1444 },
1445 )
1446 .await;
1447 }
1448
1449 #[fuchsia::test]
1450 async fn get_already_cached_missing_content_blob() {
1451 let (client, mut server) = MockPackageCache::new();
1452
1453 let ((), ()) = future::join(
1454 async {
1455 server
1456 .expect_get(blob_info(2), fpkg::GcProtection::OpenPackageTracking)
1457 .await
1458 .expect_open_meta_blob(Ok(None))
1459 .await
1460 .expect_get_missing_blobs_client_closes_channel(vec![vec![BlobInfo {
1461 blob_id: [0; 32].into(),
1462 length: 0,
1463 }]])
1464 .await;
1465 },
1466 async move {
1467 assert_matches!(
1468 client.get_already_cached(blob_id(2)).await,
1469 Err(GetAlreadyCachedError::MissingContentBlobs(v))
1470 if v == vec![BlobInfo {
1471 blob_id: [0; 32].into(),
1472 length: 0,
1473 }]
1474 );
1475 },
1476 )
1477 .await;
1478 }
1479}