1use std::{
6 collections::HashMap,
7 io,
8 path::{Path, PathBuf},
9 pin::Pin,
10 sync::Arc,
11 task::{Context, Poll},
12 time::Instant,
13};
14
15use oci_client::{Client, manifest::ImageIndexEntry};
16use tokio::{
17 io::{AsyncRead, ReadBuf},
18 sync::Semaphore,
19 task::JoinHandle,
20};
21
22use crate::{
23 cache::{
24 self, CachedImageMetadata, CachedLayerMetadata, GlobalCache, VmdkWriteMode,
25 lock::{lock_exclusive, open_lock_file},
26 },
27 config::ImageConfig,
28 digest::Digest,
29 erofs,
30 error::{ImageError, ImageResult},
31 layer::Layer,
32 platform::Platform,
33 progress::{self, PullProgress, PullProgressHandle, PullProgressSender},
34 pull::{PullOptions, PullPolicy, PullResult, RootfsMaterialization},
35 tar::{self, Compression},
36 tree::{
37 DeviceNode, DirectoryNode, FileData, FileTree, RegularFileId, RegularFileNode,
38 ResourceLimits, SymlinkNode, TreeNode,
39 },
40};
41
42use super::{RegistryBuilder, manifest::OciManifest};
43
44const MATERIALIZE_PROGRESS_EMIT_BYTES: u64 = 256 * 1024;
50
51const MAX_LAYER_PIPELINE_CONCURRENCY: usize = 16;
53
54pub struct Registry {
60 pub(super) client: Client,
61 pub(super) auth: oci_client::secrets::RegistryAuth,
62 pub(super) platform: Platform,
63 pub(super) cache: GlobalCache,
64}
65
66#[derive(Debug, Clone)]
68struct LayerDescriptor {
69 digest: Digest,
70 media_type: Option<String>,
71 size: Option<u64>,
72}
73
74struct CachedPullInfo {
75 result: PullResult,
76 metadata: CachedImageMetadata,
77}
78
79struct LayerPipelineFailure {
80 error: ImageError,
81}
82
83struct MaterializeLayersRequest<'a> {
84 oci_ref: &'a oci_client::Reference,
85 manifest_digest: &'a Digest,
86 layer_descriptors: &'a [LayerDescriptor],
87 diff_ids: &'a [String],
88 force: bool,
89 materialization: RootfsMaterialization,
90 progress: Option<PullProgressSender>,
91 staged_layers: Option<Arc<crate::archive::StagedLayerGuard>>,
92}
93
94struct LayerPipelineTreeSuccess {
97 layer_index: usize,
98 tree: Option<FileTree>,
99 data_map: Option<erofs::ErofsDataMap>,
100}
101
102struct MaterializeProgressReader<R> {
109 inner: R,
110 progress: Option<PullProgressSender>,
111 layer_index: usize,
112 total_bytes: u64,
113 bytes_read: u64,
114 last_emitted_bytes: u64,
115}
116
117impl<R> MaterializeProgressReader<R> {
122 fn new(
123 inner: R,
124 progress: Option<PullProgressSender>,
125 layer_index: usize,
126 total_bytes: u64,
127 ) -> Self {
128 Self {
129 inner,
130 progress,
131 layer_index,
132 total_bytes: total_bytes.max(1),
133 bytes_read: 0,
134 last_emitted_bytes: 0,
135 }
136 }
137}
138
139impl Registry {
140 pub fn new(platform: Platform, cache: GlobalCache) -> ImageResult<Self> {
142 Self::builder(platform, cache).build()
143 }
144
145 pub fn builder(platform: Platform, cache: GlobalCache) -> RegistryBuilder {
147 RegistryBuilder::new(platform, cache)
148 }
149
150 pub fn pull_cached(
152 cache: &GlobalCache,
153 reference: &oci_client::Reference,
154 options: &PullOptions,
155 ) -> ImageResult<Option<(PullResult, CachedImageMetadata)>> {
156 Ok(resolve_cached_pull_result(cache, reference, options)?
157 .map(|cached| (cached.result, cached.metadata)))
158 }
159
160 pub async fn pull_cached_async(
162 cache: &GlobalCache,
163 reference: &oci_client::Reference,
164 options: &PullOptions,
165 ) -> ImageResult<Option<(PullResult, CachedImageMetadata)>> {
166 let cache = cache.clone();
167 let reference = reference.clone();
168 let options = options.clone();
169 tokio::task::spawn_blocking(move || Self::pull_cached(&cache, &reference, &options))
170 .await
171 .map_err(|error| ImageError::Io(io::Error::other(error)))?
172 }
173
174 pub async fn pull_cached_by_manifest_digest(
180 cache: &GlobalCache,
181 manifest_digest: &Digest,
182 ) -> ImageResult<Option<(PullResult, CachedImageMetadata)>> {
183 Ok(
184 resolve_cached_pull_result_by_manifest_digest_async(cache, manifest_digest, false)
185 .await?
186 .map(|cached| (cached.result, cached.metadata)),
187 )
188 }
189
190 pub async fn pull_snapshot_cached(
193 cache: &GlobalCache,
194 references: &[oci_client::Reference],
195 manifest_digest: &Digest,
196 materialization: RootfsMaterialization,
197 ) -> ImageResult<Option<(PullResult, CachedImageMetadata)>> {
198 let metadata_only = materialization == RootfsMaterialization::Flat;
199 let expected = manifest_digest.to_string();
200 for reference in references {
203 if let Some(metadata) = cache.read_image_metadata_async(reference).await?
204 && metadata.manifest_digest == expected
205 && let Some(cached) =
206 resolve_snapshot_metadata(cache, metadata, metadata_only).await?
207 {
208 return Ok(Some((cached.result, cached.metadata)));
209 }
210 }
211 Ok(resolve_cached_pull_result_by_manifest_digest_async(
212 cache,
213 manifest_digest,
214 metadata_only,
215 )
216 .await?
217 .map(|cached| (cached.result, cached.metadata)))
218 }
219
220 pub async fn pull_snapshot_metadata(
223 &self,
224 reference: &oci_client::Reference,
225 ) -> ImageResult<PullResult> {
226 let expected = reference.digest().ok_or_else(|| {
227 ImageError::ManifestParse("snapshot metadata requires a digest-pinned reference".into())
228 })?;
229 let (manifest_bytes, digest, config_bytes) =
230 self.fetch_manifest_and_config(reference).await?;
231 if digest != expected {
232 return Err(ImageError::ManifestParse(
233 "snapshot manifest digest differs from pinned reference".into(),
234 ));
235 }
236 let (manifest, config_bytes, resolved) = self
237 .parse_and_resolve_manifest(&manifest_bytes, config_bytes, reference)
238 .await?;
239 let (config, diff_ids) = ImageConfig::parse(&config_bytes)?;
240 let layers = self.extract_layer_digests(&manifest)?;
241 if layers.len() != diff_ids.len() {
242 return Err(ImageError::ManifestParse(
243 "snapshot manifest/config layer count mismatch".into(),
244 ));
245 }
246 let metadata = CachedImageMetadata {
247 manifest_digest: digest,
248 config_digest: manifest.config_digest().unwrap_or_default(),
249 raw_manifest_json: json_bytes_to_string(&resolved, "resolved manifest")?,
250 raw_config_json: json_bytes_to_string(&config_bytes, "image config")?,
251 config,
252 layers: layers
253 .iter()
254 .zip(diff_ids)
255 .map(|(layer, diff_id)| CachedLayerMetadata {
256 digest: layer.digest.to_string(),
257 media_type: layer.media_type.clone(),
258 size_bytes: layer.size,
259 diff_id,
260 })
261 .collect(),
262 };
263 let mut result = cached_pull_result(&metadata)?;
264 self.cache
265 .write_image_metadata_async(reference, &metadata)
266 .await?;
267 result.cached = false;
268 Ok(result)
269 }
270
271 pub async fn pull(
273 &self,
274 reference: &oci_client::Reference,
275 options: &PullOptions,
276 ) -> ImageResult<PullResult> {
277 self.pull_inner(reference, options, None).await
278 }
279
280 pub async fn materialize_cached_layers(
286 &self,
287 reference: &oci_client::Reference,
288 metadata: &CachedImageMetadata,
289 force: bool,
290 ) -> ImageResult<PullResult> {
291 self.materialize_cached_layers_inner(reference, metadata, force, None, None)
292 .await
293 }
294
295 pub async fn materialize_flat_rootfs(
297 &self,
298 manifest_digest: &Digest,
299 layer_diff_ids: &[Digest],
300 force: bool,
301 ) -> ImageResult<crate::FlatRootfsRef> {
302 let cache = self.cache.clone();
303 let platform = self.platform.clone();
304 let manifest_digest = manifest_digest.clone();
305 let layer_diff_ids = layer_diff_ids.to_vec();
306 tokio::task::spawn_blocking(move || {
307 let _leases = cache.lease_paths(
308 layer_diff_ids
309 .iter()
310 .map(|id| cache.layer_erofs_path(id))
311 .collect(),
312 )?;
313 crate::flat::materialize_flat_rootfs(
314 &cache,
315 &manifest_digest,
316 &layer_diff_ids,
317 &platform,
318 force,
319 )
320 })
321 .await
322 .map_err(|error| ImageError::Io(io::Error::other(error)))?
323 }
324
325 pub(crate) async fn materialize_cached_layers_from_paths(
326 &self,
327 reference: &oci_client::Reference,
328 metadata: &CachedImageMetadata,
329 force: bool,
330 staged_layers: Arc<crate::archive::StagedLayerGuard>,
331 progress: Option<PullProgressSender>,
332 ) -> ImageResult<PullResult> {
333 self.materialize_cached_layers_inner(
334 reference,
335 metadata,
336 force,
337 Some(staged_layers),
338 progress,
339 )
340 .await
341 }
342
343 async fn materialize_cached_layers_inner(
344 &self,
345 reference: &oci_client::Reference,
346 metadata: &CachedImageMetadata,
347 force: bool,
348 staged_layers: Option<Arc<crate::archive::StagedLayerGuard>>,
349 progress: Option<PullProgressSender>,
350 ) -> ImageResult<PullResult> {
351 let manifest_digest: Digest = metadata.manifest_digest.parse()?;
352 let layer_descriptors = metadata
353 .layers
354 .iter()
355 .map(|layer| {
356 Ok(LayerDescriptor {
357 digest: layer.digest.parse()?,
358 media_type: layer.media_type.clone(),
359 size: layer.size_bytes,
360 })
361 })
362 .collect::<ImageResult<Vec<_>>>()?;
363 let diff_ids = metadata
364 .layers
365 .iter()
366 .map(|layer| layer.diff_id.clone())
367 .collect::<Vec<_>>();
368
369 self.materialize_layer_stage(MaterializeLayersRequest {
370 oci_ref: reference,
371 manifest_digest: &manifest_digest,
372 layer_descriptors: &layer_descriptors,
373 diff_ids: &diff_ids,
374 force,
375 materialization: RootfsMaterialization::Layered,
376 progress,
377 staged_layers,
378 })
379 .await?;
380
381 let layer_diff_ids = diff_ids
382 .iter()
383 .map(|diff_id| diff_id.parse())
384 .collect::<ImageResult<Vec<Digest>>>()?;
385
386 Ok(PullResult {
387 layer_diff_ids,
388 config: metadata.config.clone(),
389 manifest_digest,
390 cached: false,
391 })
392 }
393
394 pub fn pull_with_progress(
399 &self,
400 reference: &oci_client::Reference,
401 options: &PullOptions,
402 ) -> (PullProgressHandle, JoinHandle<ImageResult<PullResult>>)
403 where
404 Self: Send + Sync + 'static,
405 {
406 let (handle, sender) = progress::progress_channel();
407 let task = self.spawn_pull_task(reference, options, sender);
408 (handle, task)
409 }
410
411 pub fn pull_with_sender(
417 &self,
418 reference: &oci_client::Reference,
419 options: &PullOptions,
420 sender: PullProgressSender,
421 ) -> JoinHandle<ImageResult<PullResult>>
422 where
423 Self: Send + Sync + 'static,
424 {
425 self.spawn_pull_task(reference, options, sender)
426 }
427
428 fn spawn_pull_task(
430 &self,
431 reference: &oci_client::Reference,
432 options: &PullOptions,
433 sender: PullProgressSender,
434 ) -> JoinHandle<ImageResult<PullResult>>
435 where
436 Self: Send + Sync + 'static,
437 {
438 let reference = reference.clone();
439 let options = options.clone();
440 let client = self.client.clone();
441 let auth = self.auth.clone();
442 let platform = self.platform.clone();
443
444 let cache = self.cache.clone();
445
446 tokio::spawn(async move {
447 let registry = Self {
448 client,
449 auth,
450 platform,
451 cache,
452 };
453 registry
454 .pull_inner(&reference, &options, Some(sender))
455 .await
456 })
457 }
458
459 async fn pull_inner(
461 &self,
462 reference: &oci_client::Reference,
463 options: &PullOptions,
464 progress: Option<PullProgressSender>,
465 ) -> ImageResult<PullResult> {
466 let scoped = Self {
467 client: self.client.clone(),
468 auth: self.auth.clone(),
469 platform: self.platform.clone(),
470 cache: self.cache.operation_or_new(),
471 };
472 scoped.pull_scoped_inner(reference, options, progress).await
473 }
474
475 async fn pull_scoped_inner(
476 &self,
477 reference: &oci_client::Reference,
478 options: &PullOptions,
479 progress: Option<PullProgressSender>,
480 ) -> ImageResult<PullResult> {
481 let _reference_lease = self
482 .cache
483 .lease_paths_async(vec![self.cache.image_metadata_path(reference)])
484 .await?;
485 let pull_started_at = Instant::now();
486 let ref_str: Arc<str> = reference.to_string().into();
487 let oci_ref = reference;
488 let image_lock_path = self.cache.image_lock_path(reference);
489 let image_lock_file = open_lock_file(&image_lock_path)?;
490 let image_lock_file = tokio::task::spawn_blocking(move || {
491 lock_exclusive(&image_lock_file)?;
492 Ok::<_, ImageError>(image_lock_file)
493 })
494 .await
495 .map_err(|e| ImageError::Io(io::Error::other(e)))??;
496 let _image_lock_guard = image_lock_file;
499
500 if let Some(cached) =
502 resolve_cached_pull_result_async(&self.cache, reference, options, &self.platform)
503 .await?
504 {
505 tracing::debug!(
506 reference = %reference,
507 elapsed_ms = pull_started_at.elapsed().as_millis(),
508 "pull resolved entirely from cached image metadata"
509 );
510
511 if let Some(ref p) = progress {
512 p.send(PullProgress::Resolving {
513 reference: ref_str.clone(),
514 });
515 p.send(PullProgress::Resolved {
516 reference: ref_str.clone(),
517 manifest_digest: cached.metadata.manifest_digest.clone().into(),
518 layer_count: cached.metadata.layers.len(),
519 total_download_bytes: cached
520 .metadata
521 .layers
522 .iter()
523 .filter_map(|layer| layer.size_bytes)
524 .reduce(|a, b| a + b),
525 });
526 p.send(PullProgress::Complete {
527 reference: ref_str,
528 layer_count: cached.metadata.layers.len(),
529 });
530 }
531
532 return Ok(cached.result);
533 }
534
535 if options.pull_policy == PullPolicy::Never {
536 return Err(ImageError::NotCached {
537 reference: reference.to_string(),
538 });
539 }
540
541 if let Some(ref p) = progress {
543 p.send(PullProgress::Resolving {
544 reference: ref_str.clone(),
545 });
546 }
547
548 let resolve_started_at = Instant::now();
549 let (manifest_bytes, manifest_digest, config_bytes) =
550 self.fetch_manifest_and_config(oci_ref).await?;
551
552 let manifest_digest: Digest = manifest_digest.parse()?;
553
554 let (manifest, config_bytes, resolved_manifest_bytes) = self
557 .parse_and_resolve_manifest(&manifest_bytes, config_bytes, oci_ref)
558 .await?;
559
560 let (image_config, diff_ids) = ImageConfig::parse(&config_bytes)?;
562
563 let layer_descriptors = self.extract_layer_digests(&manifest)?;
565
566 if diff_ids.len() != layer_descriptors.len() {
568 return Err(ImageError::ManifestParse(format!(
569 "layer count mismatch: config has {} diff_ids but manifest has {} layers",
570 diff_ids.len(),
571 layer_descriptors.len()
572 )));
573 }
574
575 let layer_count = layer_descriptors.len();
576 let total_bytes: Option<u64> = {
577 let sum: u64 = layer_descriptors
578 .iter()
579 .filter_map(|layer| layer.size)
580 .sum();
581 if sum > 0 { Some(sum) } else { None }
582 };
583
584 tracing::debug!(
585 reference = %reference,
586 layer_count,
587 elapsed_ms = resolve_started_at.elapsed().as_millis(),
588 "pull resolved manifest and layer descriptors"
589 );
590
591 if let Some(ref p) = progress {
592 p.send(PullProgress::Resolved {
593 reference: ref_str.clone(),
594 manifest_digest: manifest_digest.to_string().into(),
595 layer_count,
596 total_download_bytes: total_bytes,
597 });
598 }
599
600 tokio::task::yield_now().await;
603
604 {
606 let mut seen = std::collections::HashSet::new();
607 for desc in &layer_descriptors {
608 if !seen.insert(&desc.digest) {
609 tracing::warn!(
610 digest = %desc.digest,
611 "manifest contains duplicate layer digest; \
612 per-layer processing will be serialized for this digest"
613 );
614 }
615 }
616 }
617
618 self.materialize_layer_stage(MaterializeLayersRequest {
620 oci_ref,
621 manifest_digest: &manifest_digest,
622 layer_descriptors: &layer_descriptors,
623 diff_ids: &diff_ids,
624 force: options.force,
625 materialization: options.materialization,
626 progress: progress.clone(),
627 staged_layers: None,
628 })
629 .await?;
630
631 let layer_diff_ids: Vec<Digest> = diff_ids
632 .iter()
633 .map(|diff_id| diff_id.parse())
634 .collect::<ImageResult<Vec<Digest>>>()?;
635
636 if options.materialization.includes_flat() {
637 self.materialize_flat_rootfs(&manifest_digest, &layer_diff_ids, options.force)
638 .await?;
639 }
640
641 for layer_desc in &layer_descriptors {
644 let layer = Layer::new(layer_desc.digest.clone(), &self.cache);
645 let _ = tokio::fs::remove_file(&layer.tar_path_ref()).await;
646 }
647
648 let cached_image = CachedImageMetadata {
650 manifest_digest: manifest_digest.to_string(),
651 config_digest: manifest.config_digest().unwrap_or_default(),
652 raw_manifest_json: json_bytes_to_string(&resolved_manifest_bytes, "resolved manifest")?,
653 raw_config_json: json_bytes_to_string(&config_bytes, "image config")?,
654 config: image_config.clone(),
655 layers: layer_descriptors
656 .iter()
657 .enumerate()
658 .map(|(i, layer)| CachedLayerMetadata {
659 digest: layer.digest.to_string(),
660 media_type: layer.media_type.clone(),
661 size_bytes: layer.size,
662 diff_id: diff_ids.get(i).cloned().unwrap_or_default(),
663 })
664 .collect(),
665 };
666 self.cache
667 .write_image_metadata_async(reference, &cached_image)
668 .await?;
669
670 tracing::debug!(
671 reference = %reference,
672 layer_count,
673 elapsed_ms = pull_started_at.elapsed().as_millis(),
674 "pull completed and cached image metadata was persisted"
675 );
676
677 if let Some(ref p) = progress {
678 p.send(PullProgress::Complete {
679 reference: ref_str,
680 layer_count,
681 });
682 }
683
684 Ok(PullResult {
685 layer_diff_ids,
686 config: image_config,
687 manifest_digest,
688 cached: false,
689 })
690 }
691
692 async fn fetch_manifest_and_config(
694 &self,
695 reference: &oci_client::Reference,
696 ) -> ImageResult<(Vec<u8>, String, Vec<u8>)> {
697 let (manifest, manifest_digest, config) = self
698 .client
699 .pull_manifest_and_config(reference, &self.auth)
700 .await?;
701
702 let manifest_bytes = serde_json::to_vec(&manifest)
703 .map_err(|e| ImageError::ManifestParse(format!("failed to serialize manifest: {e}")))?;
704
705 Ok((manifest_bytes, manifest_digest, config.into_bytes()))
706 }
707
708 async fn parse_and_resolve_manifest(
714 &self,
715 manifest_bytes: &[u8],
716 config_bytes: Vec<u8>,
717 reference: &oci_client::Reference,
718 ) -> ImageResult<(OciManifest, Vec<u8>, Vec<u8>)> {
719 let media_type = detect_manifest_media_type(manifest_bytes);
721
722 let manifest = OciManifest::parse(manifest_bytes, &media_type)?;
723
724 if manifest.is_index() {
725 self.resolve_platform_manifest(manifest_bytes, reference)
727 .await
728 } else {
729 Ok((manifest, config_bytes, manifest_bytes.to_vec()))
730 }
731 }
732
733 async fn resolve_platform_manifest(
737 &self,
738 index_bytes: &[u8],
739 reference: &oci_client::Reference,
740 ) -> ImageResult<(OciManifest, Vec<u8>, Vec<u8>)> {
741 let index: oci_spec::image::ImageIndex = serde_json::from_slice(index_bytes)
742 .map_err(|e| ImageError::ManifestParse(format!("failed to parse index: {e}")))?;
743
744 let manifests = index.manifests();
745
746 let mut best_match: Option<&oci_spec::image::Descriptor> = None;
748 let mut exact_variant = false;
749
750 for entry in manifests {
751 if entry.media_type().to_string().contains("attestation") {
753 continue;
754 }
755
756 let platform = match entry.platform().as_ref() {
757 Some(p) => p,
758 None => continue,
759 };
760
761 if *platform.os() != self.platform.os {
763 continue;
764 }
765
766 if *platform.architecture() != self.platform.arch {
768 continue;
769 }
770
771 if let Some(ref target_variant) = self.platform.variant {
773 if let Some(entry_variant) = platform.variant().as_ref()
774 && entry_variant == target_variant
775 {
776 best_match = Some(entry);
777 exact_variant = true;
778 continue;
779 }
780 if !exact_variant {
781 best_match = Some(entry);
782 }
783 } else {
784 best_match = Some(entry);
785 }
786 }
787
788 let entry = best_match.ok_or_else(|| ImageError::PlatformNotFound {
789 reference: reference.to_string(),
790 os: self.platform.os.clone(),
791 arch: self.platform.arch.clone(),
792 })?;
793
794 let digest = entry.digest();
795
796 let platform_ref = format!(
798 "{}/{}@{}",
799 reference.registry(),
800 reference.repository(),
801 digest
802 );
803 let platform_ref: oci_client::Reference = platform_ref.parse().map_err(|e| {
804 ImageError::ManifestParse(format!("failed to parse platform reference: {e}"))
805 })?;
806
807 let (manifest_bytes, _digest, config_bytes) =
808 self.fetch_manifest_and_config(&platform_ref).await?;
809
810 let media_type = detect_manifest_media_type(&manifest_bytes);
811 let manifest = OciManifest::parse(&manifest_bytes, &media_type)?;
812 Ok((manifest, config_bytes, manifest_bytes))
813 }
814
815 fn extract_layer_digests(&self, manifest: &OciManifest) -> ImageResult<Vec<LayerDescriptor>> {
817 match manifest {
818 OciManifest::Image(m) => {
819 let layers: Vec<LayerDescriptor> = m
820 .layers()
821 .iter()
822 .map(|desc| {
823 let digest: Digest = desc.digest().to_string().parse().map_err(|_| {
824 ImageError::ManifestParse(format!(
825 "invalid layer digest: {}",
826 desc.digest()
827 ))
828 })?;
829 let size = if desc.size() > 0 {
830 Some(desc.size())
831 } else {
832 None
833 };
834 Ok(LayerDescriptor {
835 digest,
836 media_type: Some(desc.media_type().to_string()),
837 size,
838 })
839 })
840 .collect::<ImageResult<Vec<_>>>()?;
841 Ok(layers)
842 }
843 OciManifest::Index(_) => Err(ImageError::ManifestParse(
844 "cannot extract layers from an index — resolve platform first".to_string(),
845 )),
846 }
847 }
848
849 async fn materialize_layer_stage(
851 &self,
852 request: MaterializeLayersRequest<'_>,
853 ) -> ImageResult<()> {
854 let MaterializeLayersRequest {
855 oci_ref,
856 manifest_digest,
857 layer_descriptors,
858 diff_ids,
859 force,
860 materialization,
861 progress,
862 staged_layers,
863 } = request;
864
865 let validated_diff_ids: Vec<Digest> = diff_ids
868 .iter()
869 .enumerate()
870 .map(|(i, id)| {
871 id.parse::<Digest>().map_err(|_| {
872 ImageError::ManifestParse(format!("invalid diff_id at layer {i}: {id}"))
873 })
874 })
875 .collect::<ImageResult<Vec<_>>>()?;
876
877 let mut protected = vec![
878 self.cache.fsmeta_erofs_path(manifest_digest),
879 self.cache.vmdk_path(manifest_digest),
880 ];
881 protected.extend(
882 validated_diff_ids
883 .iter()
884 .map(|id| self.cache.layer_erofs_path(id)),
885 );
886 let _leases = self.cache.lease_paths_async(protected).await?;
887
888 let fsmeta_path = self.cache.fsmeta_erofs_path(manifest_digest);
897 let vmdk_path = self.cache.vmdk_path(manifest_digest);
898 let fsmeta_valid = cache::is_valid_erofs_artifact_async(&fsmeta_path).await;
899 let vmdk_valid = path_exists_async(&vmdk_path).await;
900 let all_layers_valid =
901 all_layers_materialized_async(&self.cache, &validated_diff_ids).await;
902
903 if all_layers_valid
904 && (!materialization.includes_layered() || (fsmeta_valid && vmdk_valid))
905 && !force
906 {
907 return Ok(());
908 }
909
910 if materialization.includes_layered()
911 && all_layers_valid
912 && fsmeta_valid
913 && !vmdk_valid
914 && !force
915 {
916 return self
917 .regenerate_vmdk_only(manifest_digest, &validated_diff_ids, progress.as_ref())
918 .await;
919 }
920
921 let layer_force = force;
926 let layer_concurrency = layer_pipeline_concurrency(layer_descriptors.len());
927 let semaphore = Arc::new(Semaphore::new(layer_concurrency));
928
929 let layer_tasks: Vec<_> = layer_descriptors
930 .iter()
931 .enumerate()
932 .map(|(i, layer_desc)| {
933 let layer = Layer::new(layer_desc.digest.clone(), &self.cache);
934 let client = self.client.clone();
935 let oci_ref = oci_ref.clone();
936 let size = layer_desc.size;
937 let progress = progress.clone();
938 let media_type = layer_desc.media_type.clone();
939 let diff_id = diff_ids[i].clone();
940 let staged_tar_path = staged_layers
941 .as_ref()
942 .and_then(|layers| layers.get(&layer_desc.digest.to_string()).cloned());
943
944 let diff_id_digest: Digest = validated_diff_ids[i].clone();
945 let erofs_path = self.cache.layer_erofs_path(&diff_id_digest);
946 let lock_path = self.cache.layer_erofs_lock_path(&diff_id_digest);
947 let tmp_dir = self.cache.tmp_dir().to_path_buf();
948 let semaphore = Arc::clone(&semaphore);
949 let leases = _leases.clone();
951 let staging = staged_layers.clone();
952 tokio::spawn(async move {
953 let _leases = leases;
954 let _staging = staging;
955 let _permit =
956 semaphore
957 .acquire_owned()
958 .await
959 .map_err(|e| LayerPipelineFailure {
960 error: ImageError::Io(io::Error::other(format!(
961 "layer pipeline semaphore closed: {e}"
962 ))),
963 })?;
964 let layer_started_at = Instant::now();
965
966 if cache::is_valid_erofs_artifact_async(&erofs_path).await && !layer_force {
967 if let Some(ref p) = progress {
968 p.send(PullProgress::LayerMaterializeComplete {
969 layer_index: i,
970 diff_id: diff_id.clone().into(),
971 });
972 }
973
974 tracing::debug!(
975 layer_index = i,
976 diff_id = %diff_id,
977 elapsed_ms = layer_started_at.elapsed().as_millis(),
978 "layer reused existing EROFS image"
979 );
980
981 return Ok::<_, LayerPipelineFailure>(LayerPipelineTreeSuccess {
982 layer_index: i,
983 tree: None,
984 data_map: None,
985 });
986 }
987
988 if staged_tar_path.is_none() && _staging.is_some() {
989 return Err(LayerPipelineFailure {
990 error: ImageError::ManifestParse(format!(
991 "archive has no staged layer {}",
992 diff_id
993 )),
994 });
995 }
996 if staged_tar_path.is_none()
997 && let Err(error) = layer
998 .download(&client, &oci_ref, size, force, progress.as_ref(), i)
999 .await
1000 {
1001 return Err(LayerPipelineFailure { error });
1002 }
1003
1004 let lock_file = open_lock_file(&lock_path)
1006 .map_err(|e| LayerPipelineFailure { error: e })?;
1007 let lock_file = tokio::task::spawn_blocking(move || {
1008 lock_exclusive(&lock_file)?;
1009 Ok::<_, ImageError>(lock_file)
1010 })
1011 .await
1012 .map_err(|e| LayerPipelineFailure {
1013 error: ImageError::Io(io::Error::other(e)),
1014 })?
1015 .map_err(|e| LayerPipelineFailure { error: e })?;
1016 let materializer_guard = lock_file;
1017
1018 if cache::is_valid_erofs_artifact_async(&erofs_path).await && !layer_force {
1020 if let Some(ref p) = progress {
1021 p.send(PullProgress::LayerMaterializeComplete {
1022 layer_index: i,
1023 diff_id: diff_id.clone().into(),
1024 });
1025 }
1026 return Ok::<_, LayerPipelineFailure>(LayerPipelineTreeSuccess {
1027 layer_index: i,
1028 tree: None,
1029 data_map: None,
1030 });
1031 }
1032
1033 if let Some(ref p) = progress {
1034 p.send(PullProgress::LayerMaterializeStarted {
1035 layer_index: i,
1036 diff_id: diff_id.clone().into(),
1037 });
1038 }
1039
1040 let tar_path = staged_tar_path
1041 .clone()
1042 .unwrap_or_else(|| layer.tar_path_ref());
1043 let tar_size =
1044 tokio::fs::metadata(&tar_path)
1045 .await
1046 .map_err(|e| LayerPipelineFailure {
1047 error: ImageError::Cache {
1048 path: tar_path.clone(),
1049 source: e,
1050 },
1051 })?;
1052 let tar_file = tokio::fs::File::open(&tar_path).await.map_err(|e| {
1053 LayerPipelineFailure {
1054 error: ImageError::Cache {
1055 path: tar_path.clone(),
1056 source: e,
1057 },
1058 }
1059 })?;
1060
1061 let compression =
1062 Compression::from_media_type(media_type.as_deref().unwrap_or(""));
1063 let limits = ResourceLimits::default();
1064 let spool_path = layer_work_path(&tmp_dir, &diff_id_digest, "spool");
1065 let ingest_started_at = Instant::now();
1066 let ingest_result = tar::ingest_compressed_tar(
1067 MaterializeProgressReader::new(
1068 tar_file,
1069 progress.clone(),
1070 i,
1071 tar_size.len(),
1072 ),
1073 compression,
1074 &limits,
1075 Some(&spool_path),
1076 )
1077 .await
1078 .map_err(|e| LayerPipelineFailure {
1079 error: ImageError::Materialize {
1080 digest: diff_id.clone(),
1081 message: format!("tar ingestion failed: {e}"),
1082 source: None,
1083 },
1084 })?;
1085
1086 let expected_diff_hex = diff_id_digest.hex();
1090 if ingest_result.uncompressed_digest != expected_diff_hex {
1091 return Err(LayerPipelineFailure {
1092 error: ImageError::DigestMismatch {
1093 digest: diff_id.clone(),
1094 expected: format!("sha256:{expected_diff_hex}"),
1095 actual: format!("sha256:{}", ingest_result.uncompressed_digest),
1096 },
1097 });
1098 }
1099 let tree = ingest_result.tree;
1100
1101 tracing::debug!(
1102 layer_index = i,
1103 diff_id = %diff_id,
1104 tar_bytes = tar_size.len(),
1105 elapsed_ms = ingest_started_at.elapsed().as_millis(),
1106 "layer tar ingestion completed (diff_id verified)"
1107 );
1108
1109 if let Some(ref p) = progress {
1110 p.send(PullProgress::LayerMaterializeWriting { layer_index: i });
1111 }
1112
1113 let temp_path = layer_work_path(&tmp_dir, &diff_id_digest, "erofs.part");
1116 let erofs_final = erofs_path.clone();
1117 let diff_id_for_join = diff_id.clone();
1118 let write_started_at = Instant::now();
1119 let worker_leases = _leases.clone();
1120 let (data_map, mut tree) = tokio::task::spawn_blocking(move || {
1121 let _protection = (materializer_guard, worker_leases);
1122 let data_map = erofs::write_erofs(&tree, &temp_path)?;
1123 std::fs::rename(&temp_path, &erofs_final).map_err(erofs::ErofsError::Io)?;
1124 Ok::<(erofs::ErofsDataMap, FileTree), erofs::ErofsError>((data_map, tree))
1125 })
1126 .await
1127 .map_err(|e| LayerPipelineFailure {
1128 error: ImageError::Materialize {
1129 digest: diff_id_for_join.clone(),
1130 message: format!("EROFS write task failed: {e}"),
1131 source: None,
1132 },
1133 })?
1134 .map_err(|e| LayerPipelineFailure {
1135 error: ImageError::Materialize {
1136 digest: diff_id.clone(),
1137 message: format!("EROFS write failed: {e}"),
1138 source: None,
1139 },
1140 })?;
1141
1142 tree.strip_file_data();
1145
1146 tracing::debug!(
1147 layer_index = i,
1148 diff_id = %diff_id,
1149 elapsed_ms = write_started_at.elapsed().as_millis(),
1150 total_elapsed_ms = layer_started_at.elapsed().as_millis(),
1151 "layer EROFS image write completed"
1152 );
1153
1154 let _ = tokio::fs::remove_file(&spool_path).await;
1158
1159 if let Some(ref p) = progress {
1160 p.send(PullProgress::LayerMaterializeComplete {
1161 layer_index: i,
1162 diff_id: diff_id.clone().into(),
1163 });
1164 }
1165
1166 Ok::<_, LayerPipelineFailure>(LayerPipelineTreeSuccess {
1167 layer_index: i,
1168 tree: Some(tree),
1169 data_map: Some(data_map),
1170 })
1171 })
1172 })
1173 .collect();
1174
1175 let mut layer_results = wait_for_layer_tree_pipeline(layer_tasks).await?;
1177 layer_results.sort_by_key(|r| r.layer_index);
1178
1179 if !materialization.includes_layered() {
1180 return Ok(());
1181 }
1182
1183 let fsmeta_path = self.cache.fsmeta_erofs_path(manifest_digest);
1185 let vmdk_path = self.cache.vmdk_path(manifest_digest);
1186
1187 if cache::is_valid_erofs_artifact_async(&fsmeta_path).await
1188 && path_exists_async(&vmdk_path).await
1189 && !force
1190 {
1191 tracing::debug!(
1192 manifest_digest = %manifest_digest,
1193 "fsmeta + VMDK already cached, skipping generation"
1194 );
1195 return Ok(());
1196 }
1197
1198 let fsmeta_lock_path = self.cache.fsmeta_erofs_lock_path(manifest_digest);
1200 let fsmeta_lock_file = open_lock_file(&fsmeta_lock_path)?;
1201 let fsmeta_lock_file = tokio::task::spawn_blocking(move || {
1202 lock_exclusive(&fsmeta_lock_file)?;
1203 Ok::<_, ImageError>(fsmeta_lock_file)
1204 })
1205 .await
1206 .map_err(|e| ImageError::Io(io::Error::other(e)))??;
1207 let materializer_guard = Arc::new(fsmeta_lock_file);
1208
1209 if cache::is_valid_erofs_artifact_async(&fsmeta_path).await
1211 && path_exists_async(&vmdk_path).await
1212 && !force
1213 {
1214 return Ok(());
1215 }
1216
1217 let cache = self.cache.clone();
1222 let diff_ids = diff_ids.to_vec();
1223 let manifest_digest_for_inputs = manifest_digest.to_string();
1224 let input_guard = materializer_guard.clone();
1225 let (layer_trees, layer_data_maps) = tokio::task::spawn_blocking(move || {
1226 let _guard = input_guard;
1227 collect_layered_inputs_from_erofs(
1228 &cache,
1229 &diff_ids,
1230 layer_results,
1231 &manifest_digest_for_inputs,
1232 )
1233 })
1234 .await
1235 .map_err(|error| ImageError::Io(io::Error::other(error)))??;
1236
1237 if let Some(ref p) = progress {
1239 p.send(PullProgress::StitchMergingTrees {
1240 layer_count: layer_trees.len(),
1241 });
1242 }
1243 let (merged_tree, provenance) = crate::tree::merge_layers_with_provenance(layer_trees);
1244
1245 let fsmeta_path_for_write = fsmeta_path.clone();
1247 let vmdk_path_for_write = vmdk_path.clone();
1248 let work_dir = self.cache.work_dir(manifest_digest);
1249 let manifest_digest_str = manifest_digest.to_string();
1250
1251 let layer_erofs_paths: Vec<std::path::PathBuf> = validated_diff_ids
1253 .iter()
1254 .map(|d| self.cache.layer_erofs_path(d))
1255 .collect();
1256
1257 let stitch_progress = progress.clone();
1258 let worker_cache = self.cache.clone();
1259 tokio::task::spawn_blocking(move || {
1260 let _protection = (materializer_guard, worker_cache);
1262 std::fs::create_dir_all(&work_dir).map_err(|e| ImageError::Cache {
1263 path: work_dir.clone(),
1264 source: e,
1265 })?;
1266 let _work_guard = scopeguard::guard((), |_| {
1267 let _ = std::fs::remove_dir_all(&work_dir);
1268 });
1269
1270 if let Some(ref p) = stitch_progress {
1272 p.send(PullProgress::StitchWritingFsmeta);
1273 }
1274 let temp_fsmeta = work_dir.join("fsmeta.erofs");
1275 erofs::fsmeta::write_fsmeta(&merged_tree, &provenance, &layer_data_maps, &temp_fsmeta)
1276 .map_err(|e| ImageError::Materialize {
1277 digest: manifest_digest_str.clone(),
1278 message: format!("fsmeta write failed: {e}"),
1279 source: None,
1280 })?;
1281
1282 std::fs::rename(&temp_fsmeta, &fsmeta_path_for_write).map_err(|e| {
1283 ImageError::Cache {
1284 path: fsmeta_path_for_write.clone(),
1285 source: e,
1286 }
1287 })?;
1288
1289 if let Some(ref p) = stitch_progress {
1291 p.send(PullProgress::StitchWritingVmdk);
1292 }
1293 let temp_vmdk = work_dir.join("rootfs.vmdk");
1294 let mut extents: Vec<&std::path::Path> = vec![&fsmeta_path_for_write];
1295 extents.extend(layer_erofs_paths.iter().map(|p| p.as_path()));
1296
1297 crate::stitch::write_vmdk_descriptor(&temp_vmdk, &extents).map_err(|e| {
1298 ImageError::Materialize {
1299 digest: manifest_digest_str.clone(),
1300 message: format!("VMDK write failed: {e}"),
1301 source: None,
1302 }
1303 })?;
1304
1305 std::fs::rename(&temp_vmdk, &vmdk_path_for_write).map_err(|e| ImageError::Cache {
1306 path: vmdk_path_for_write.clone(),
1307 source: e,
1308 })?;
1309
1310 Ok::<(), ImageError>(())
1311 })
1312 .await
1313 .map_err(|e| ImageError::Io(io::Error::other(e)))??;
1314
1315 if let Some(ref p) = progress {
1316 p.send(PullProgress::StitchComplete);
1317 }
1318
1319 Ok(())
1320 }
1321
1322 async fn regenerate_vmdk_only(
1328 &self,
1329 manifest_digest: &Digest,
1330 validated_diff_ids: &[Digest],
1331 progress: Option<&PullProgressSender>,
1332 ) -> ImageResult<()> {
1333 let cache = self.cache.clone();
1334 let manifest_digest = manifest_digest.clone();
1335 let diff_ids = validated_diff_ids.to_vec();
1336 let progress = progress.cloned();
1337
1338 tokio::task::spawn_blocking(move || {
1341 cache.write_vmdk(
1342 &manifest_digest,
1343 &diff_ids,
1344 VmdkWriteMode::KeepExisting,
1345 progress.as_ref(),
1346 )
1347 })
1348 .await
1349 .map_err(|e| ImageError::Io(io::Error::other(e)))?
1350 }
1351
1352 }
1355
1356impl<R: AsyncRead + Unpin> AsyncRead for MaterializeProgressReader<R> {
1361 fn poll_read(
1362 mut self: Pin<&mut Self>,
1363 cx: &mut Context<'_>,
1364 buf: &mut ReadBuf<'_>,
1365 ) -> Poll<io::Result<()>> {
1366 let before = buf.filled().len();
1367 match Pin::new(&mut self.inner).poll_read(cx, buf) {
1368 Poll::Ready(Ok(())) => {
1369 let bytes_read = (buf.filled().len() - before) as u64;
1370 if bytes_read > 0 {
1371 self.bytes_read += bytes_read;
1372 let should_emit_progress =
1373 self.bytes_read.saturating_sub(self.last_emitted_bytes)
1374 >= MATERIALIZE_PROGRESS_EMIT_BYTES
1375 || self.bytes_read >= self.total_bytes;
1376
1377 if should_emit_progress {
1378 if let Some(progress) = &self.progress {
1379 progress.send(PullProgress::LayerMaterializeProgress {
1380 layer_index: self.layer_index,
1381 bytes_read: self.bytes_read.min(self.total_bytes),
1382 total_bytes: self.total_bytes,
1383 });
1384 }
1385 self.last_emitted_bytes = self.bytes_read;
1386 }
1387 }
1388
1389 Poll::Ready(Ok(()))
1390 }
1391 Poll::Ready(Err(error)) => Poll::Ready(Err(error)),
1392 Poll::Pending => Poll::Pending,
1393 }
1394 }
1395}
1396
1397fn detect_manifest_media_type(bytes: &[u8]) -> String {
1403 if let Ok(v) = serde_json::from_slice::<serde_json::Value>(bytes) {
1405 if let Some(mt) = v.get("mediaType").and_then(|v| v.as_str()) {
1406 return mt.to_string();
1407 }
1408
1409 if v.get("manifests").is_some() {
1411 return "application/vnd.oci.image.index.v1+json".to_string();
1412 }
1413
1414 if v.get("layers").is_some() {
1416 return "application/vnd.oci.image.manifest.v1+json".to_string();
1417 }
1418 }
1419
1420 "application/vnd.oci.image.manifest.v1+json".to_string()
1422}
1423
1424pub(super) fn resolve_platform_digest(
1426 manifests: &[ImageIndexEntry],
1427 target: &Platform,
1428) -> Option<String> {
1429 let mut arch_only_match: Option<String> = None;
1430 let target_os = target.os.to_string();
1431 let target_arch = target.arch.to_string();
1432
1433 for entry in manifests {
1434 if entry.media_type.contains("attestation") {
1435 continue;
1436 }
1437
1438 let Some(platform) = entry.platform.as_ref() else {
1439 continue;
1440 };
1441 if platform.os.to_string() != target_os || platform.architecture.to_string() != target_arch
1442 {
1443 continue;
1444 }
1445
1446 match target.variant.as_deref() {
1447 Some(target_variant) if platform.variant.as_deref() == Some(target_variant) => {
1448 return Some(entry.digest.clone());
1449 }
1450 Some(_) => {
1451 if arch_only_match.is_none() {
1452 arch_only_match = Some(entry.digest.clone());
1453 }
1454 }
1455 None => return Some(entry.digest.clone()),
1456 }
1457 }
1458
1459 arch_only_match
1460}
1461
1462fn cached_pull_result(metadata: &CachedImageMetadata) -> ImageResult<PullResult> {
1464 let manifest_digest: Digest = metadata.manifest_digest.parse()?;
1465 let layer_diff_ids = metadata
1466 .layers
1467 .iter()
1468 .map(|layer| layer.diff_id.parse())
1469 .collect::<ImageResult<Vec<Digest>>>()?;
1470
1471 Ok(PullResult {
1472 layer_diff_ids,
1473 config: metadata.config.clone(),
1474 manifest_digest,
1475 cached: true,
1476 })
1477}
1478
1479fn collect_layered_inputs_from_erofs(
1485 cache: &GlobalCache,
1486 diff_ids: &[String],
1487 results: Vec<LayerPipelineTreeSuccess>,
1488 manifest_digest: &str,
1489) -> ImageResult<(Vec<FileTree>, Vec<erofs::ErofsDataMap>)> {
1490 let mut inputs_by_diff_id: HashMap<String, (FileTree, erofs::ErofsDataMap)> = HashMap::new();
1491 for mut result in results {
1492 match (result.tree.take(), result.data_map.take()) {
1493 (Some(tree), Some(data_map)) => {
1494 let diff_id =
1495 diff_ids
1496 .get(result.layer_index)
1497 .ok_or_else(|| ImageError::Materialize {
1498 digest: manifest_digest.to_string(),
1499 message: "layer pipeline returned an out-of-range index".into(),
1500 source: None,
1501 })?;
1502 inputs_by_diff_id
1503 .entry(diff_id.clone())
1504 .or_insert((tree, data_map));
1505 }
1506 (None, None) => {}
1507 _ => {
1508 return Err(ImageError::Materialize {
1509 digest: manifest_digest.to_string(),
1510 message: "layer pipeline returned an incomplete tree/block-map pair".into(),
1511 source: None,
1512 });
1513 }
1514 }
1515 }
1516
1517 for diff_id in diff_ids {
1518 if inputs_by_diff_id.contains_key(diff_id) {
1519 continue;
1520 }
1521 let digest: Digest = diff_id.parse().map_err(|_| {
1522 ImageError::ManifestParse(format!("invalid cached layer diff_id: {diff_id}"))
1523 })?;
1524 let path = cache.layer_erofs_path(&digest);
1525 let input = read_erofs_layer_metadata(&path).map_err(|source| ImageError::Materialize {
1526 digest: diff_id.clone(),
1527 message: "failed to reconstruct cached EROFS metadata for fsmeta".into(),
1528 source: Some(Box::new(source)),
1529 })?;
1530 inputs_by_diff_id.insert(diff_id.clone(), input);
1531 }
1532
1533 let mut trees = Vec::with_capacity(diff_ids.len());
1534 let mut data_maps = Vec::with_capacity(diff_ids.len());
1535 for diff_id in diff_ids {
1536 let (tree, data_map) =
1537 inputs_by_diff_id
1538 .get(diff_id)
1539 .ok_or_else(|| ImageError::Materialize {
1540 digest: diff_id.clone(),
1541 message: "missing reconstructed EROFS layer input".into(),
1542 source: None,
1543 })?;
1544 trees.push(tree.clone());
1545 data_maps.push(data_map.clone());
1546 }
1547
1548 Ok((trees, data_maps))
1549}
1550
1551fn read_erofs_layer_metadata(path: &Path) -> io::Result<(FileTree, erofs::ErofsDataMap)> {
1553 let file = std::fs::File::open(path)?;
1554 let total_blocks = file
1555 .metadata()?
1556 .len()
1557 .div_ceil(u64::from(crate::erofs::format::EROFS_BLKSIZ));
1558 let total_blocks = u32::try_from(total_blocks)
1559 .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, "EROFS layer is too large"))?;
1560 let mut reader = erofs::ErofsReader::new(file)?;
1561 let (root_metadata, root_xattrs) = reader.root_directory_metadata()?;
1562 let mut root = DirectoryNode::new(root_metadata);
1563 root.xattrs = root_xattrs;
1564 let mut tree = FileTree { root };
1565 let mut hardlinks: HashMap<u32, RegularFileId> = HashMap::new();
1566 let mut file_blocks = HashMap::new();
1567
1568 reader.walk_entries_with_path_bytes::<io::Error, _>(|reader, path, entry| {
1569 let node = match entry.kind {
1570 erofs::ErofsEntryKind::RegularFile => {
1571 let id = *hardlinks.entry(entry.nid).or_default();
1572 file_blocks.insert(entry.path.clone(), reader.file_block_mapping(entry.nid)?);
1573 TreeNode::RegularFile(RegularFileNode {
1574 id,
1575 metadata: entry.metadata,
1576 xattrs: entry.xattrs,
1577 data: FileData::Memory(Vec::new()),
1580 nlink: 1,
1581 })
1582 }
1583 erofs::ErofsEntryKind::Directory => {
1584 let mut directory = DirectoryNode::new(entry.metadata);
1585 directory.xattrs = entry.xattrs;
1586 TreeNode::Directory(directory)
1587 }
1588 erofs::ErofsEntryKind::Symlink => TreeNode::Symlink(SymlinkNode {
1589 metadata: entry.metadata,
1590 target: reader.read_link_by_nid(entry.nid)?,
1591 }),
1592 erofs::ErofsEntryKind::CharDevice | erofs::ErofsEntryKind::BlockDevice => {
1593 let (major, minor) = entry.rdev.ok_or_else(|| {
1594 io::Error::new(
1595 io::ErrorKind::InvalidData,
1596 "EROFS device entry is missing major/minor",
1597 )
1598 })?;
1599 let device = DeviceNode {
1600 metadata: entry.metadata,
1601 major,
1602 minor,
1603 };
1604 if entry.kind == erofs::ErofsEntryKind::CharDevice {
1605 TreeNode::CharDevice(device)
1606 } else {
1607 TreeNode::BlockDevice(device)
1608 }
1609 }
1610 erofs::ErofsEntryKind::Fifo => TreeNode::Fifo(entry.metadata),
1611 erofs::ErofsEntryKind::Socket => TreeNode::Socket(entry.metadata),
1612 };
1613 tree.insert(path, node).map_err(io::Error::other)
1614 })?;
1615
1616 Ok((
1617 tree,
1618 erofs::ErofsDataMap {
1619 file_blocks,
1620 total_blocks,
1621 },
1622 ))
1623}
1624
1625fn resolve_cached_pull_result(
1626 cache: &GlobalCache,
1627 reference: &oci_client::Reference,
1628 options: &PullOptions,
1629) -> ImageResult<Option<CachedPullInfo>> {
1630 resolve_cached_pull_result_for_platform(cache, reference, options, &Platform::host_linux())
1631}
1632
1633fn resolve_cached_pull_result_for_platform(
1634 cache: &GlobalCache,
1635 reference: &oci_client::Reference,
1636 options: &PullOptions,
1637 platform: &Platform,
1638) -> ImageResult<Option<CachedPullInfo>> {
1639 if options.force || options.pull_policy == PullPolicy::Always {
1640 return Ok(None);
1641 }
1642
1643 let Some(metadata) = cache.read_image_metadata(reference)? else {
1644 return Ok(None);
1645 };
1646
1647 let _leases = match cache.metadata_paths(&metadata) {
1648 Ok(paths) => cache.lease_paths(paths)?,
1649 Err(_) => return Ok(None),
1650 };
1651
1652 let cached_diff_ids = match metadata
1654 .layers
1655 .iter()
1656 .map(|layer| layer.diff_id.parse())
1657 .collect::<ImageResult<Vec<Digest>>>()
1658 {
1659 Ok(digests) => digests,
1660 Err(_) => return Ok(None),
1661 };
1662 if !cache.all_layers_materialized(&cached_diff_ids) {
1663 return Ok(None);
1664 }
1665
1666 let manifest_digest = match metadata.manifest_digest.parse::<Digest>() {
1667 Ok(digest) => digest,
1668 Err(_) => return Ok(None),
1669 };
1670 if options.materialization.includes_layered()
1671 && (!cache.is_fsmeta_materialized(&manifest_digest)
1672 || !cache.is_vmdk_materialized(&manifest_digest))
1673 {
1674 return Ok(None);
1675 }
1676 if options.materialization.includes_flat()
1677 && crate::flat::read_current_flat_ref(cache, &manifest_digest, &cached_diff_ids, platform)?
1678 .is_none()
1679 {
1680 return Ok(None);
1681 }
1682
1683 let result = match cached_pull_result(&metadata) {
1684 Ok(result) => result,
1685 Err(_) => return Ok(None),
1686 };
1687
1688 Ok(Some(CachedPullInfo { result, metadata }))
1689}
1690
1691async fn wait_for_layer_tree_pipeline(
1692 layer_tasks: Vec<JoinHandle<Result<LayerPipelineTreeSuccess, LayerPipelineFailure>>>,
1693) -> ImageResult<Vec<LayerPipelineTreeSuccess>> {
1694 let outcomes = futures::future::join_all(layer_tasks).await;
1695 let mut results = Vec::new();
1696 let mut first_error: Option<ImageError> = None;
1697
1698 for outcome in outcomes {
1699 match outcome {
1700 Ok(Ok(result)) => results.push(result),
1701 Ok(Err(failure)) => {
1702 if first_error.is_none() {
1703 first_error = Some(failure.error);
1704 }
1705 }
1706 Err(error) => {
1707 if first_error.is_none() {
1708 first_error = Some(ImageError::Io(io::Error::other(format!(
1709 "layer task failed: {error}"
1710 ))));
1711 }
1712 }
1713 }
1714 }
1715
1716 if let Some(error) = first_error {
1717 return Err(error);
1718 }
1719
1720 Ok(results)
1721}
1722
1723async fn resolve_cached_pull_result_async(
1724 cache: &GlobalCache,
1725 reference: &oci_client::Reference,
1726 options: &PullOptions,
1727 platform: &Platform,
1728) -> ImageResult<Option<CachedPullInfo>> {
1729 if options.force || options.pull_policy == PullPolicy::Always {
1730 return Ok(None);
1731 }
1732
1733 let Some(metadata) = cache.read_image_metadata_async(reference).await? else {
1734 return Ok(None);
1735 };
1736
1737 resolve_cached_metadata_pull_result_async(cache, metadata, options.materialization, platform)
1738 .await
1739}
1740
1741async fn resolve_cached_pull_result_by_manifest_digest_async(
1742 cache: &GlobalCache,
1743 manifest_digest: &Digest,
1744 metadata_only: bool,
1745) -> ImageResult<Option<CachedPullInfo>> {
1746 let expected = manifest_digest.to_string();
1747 let mut entries = tokio::fs::read_dir(cache.manifests_dir())
1748 .await
1749 .map_err(|e| ImageError::Cache {
1750 path: cache.manifests_dir().to_path_buf(),
1751 source: e,
1752 })?;
1753
1754 while let Some(entry) = entries.next_entry().await.map_err(|e| ImageError::Cache {
1755 path: cache.manifests_dir().to_path_buf(),
1756 source: e,
1757 })? {
1758 let path = entry.path();
1759 if path.extension().and_then(|s| s.to_str()) != Some("json") {
1760 continue;
1761 }
1762
1763 let Some(metadata) = cache
1764 .read_image_metadata_matching_async(path, expected.clone())
1765 .await?
1766 else {
1767 continue;
1768 };
1769
1770 if let Some(cached) = resolve_snapshot_metadata(cache, metadata, metadata_only).await? {
1771 return Ok(Some(cached));
1772 }
1773 }
1774
1775 Ok(None)
1776}
1777
1778async fn resolve_snapshot_metadata(
1779 cache: &GlobalCache,
1780 metadata: CachedImageMetadata,
1781 metadata_only: bool,
1782) -> ImageResult<Option<CachedPullInfo>> {
1783 if metadata_only {
1784 return Ok(cached_pull_result(&metadata)
1785 .ok()
1786 .map(|result| CachedPullInfo { result, metadata }));
1787 }
1788 resolve_cached_metadata_pull_result_async(
1789 cache,
1790 metadata,
1791 RootfsMaterialization::Layered,
1792 &Platform::host_linux(),
1793 )
1794 .await
1795}
1796
1797async fn resolve_cached_metadata_pull_result_async(
1798 cache: &GlobalCache,
1799 metadata: CachedImageMetadata,
1800 materialization: RootfsMaterialization,
1801 platform: &Platform,
1802) -> ImageResult<Option<CachedPullInfo>> {
1803 let _leases = match cache.metadata_paths(&metadata) {
1804 Ok(paths) => cache.lease_paths_async(paths).await?,
1805 Err(_) => return Ok(None),
1806 };
1807 let cached_diff_ids = match metadata
1808 .layers
1809 .iter()
1810 .map(|layer| layer.diff_id.parse())
1811 .collect::<ImageResult<Vec<Digest>>>()
1812 {
1813 Ok(digests) => digests,
1814 Err(_) => return Ok(None),
1815 };
1816 if !all_layers_materialized_async(cache, &cached_diff_ids).await {
1817 return Ok(None);
1818 }
1819
1820 let manifest_digest = match metadata.manifest_digest.parse::<Digest>() {
1821 Ok(digest) => digest,
1822 Err(_) => return Ok(None),
1823 };
1824 if materialization.includes_layered()
1825 && (!cache::is_valid_erofs_artifact_async(&cache.fsmeta_erofs_path(&manifest_digest)).await
1826 || !path_exists_async(&cache.vmdk_path(&manifest_digest)).await)
1827 {
1828 return Ok(None);
1829 }
1830 if materialization.includes_flat()
1831 && crate::flat::read_current_flat_ref(cache, &manifest_digest, &cached_diff_ids, platform)?
1832 .is_none()
1833 {
1834 return Ok(None);
1835 }
1836
1837 let result = match cached_pull_result(&metadata) {
1838 Ok(result) => result,
1839 Err(_) => return Ok(None),
1840 };
1841
1842 Ok(Some(CachedPullInfo { result, metadata }))
1843}
1844
1845async fn all_layers_materialized_async(cache: &GlobalCache, diff_ids: &[Digest]) -> bool {
1846 for diff_id in diff_ids {
1847 if !cache::is_valid_erofs_artifact_async(&cache.layer_erofs_path(diff_id)).await {
1848 return false;
1849 }
1850 }
1851
1852 true
1853}
1854
1855async fn path_exists_async(path: &Path) -> bool {
1856 tokio::fs::metadata(path).await.is_ok()
1857}
1858
1859fn layer_pipeline_concurrency(layer_count: usize) -> usize {
1860 let host_limit = std::thread::available_parallelism()
1861 .map(|n| n.get().saturating_mul(2))
1862 .unwrap_or(8)
1863 .clamp(4, MAX_LAYER_PIPELINE_CONCURRENCY);
1864
1865 host_limit.min(layer_count.max(1))
1866}
1867
1868fn layer_work_path(tmp_dir: &Path, diff_id: &Digest, suffix: &str) -> PathBuf {
1869 tmp_dir.join(format!("{}.{}", diff_id.to_path_safe(), suffix))
1870}
1871
1872fn json_bytes_to_string(bytes: &[u8], context: &str) -> ImageResult<String> {
1873 std::str::from_utf8(bytes)
1874 .map(str::to_owned)
1875 .map_err(|e| ImageError::ConfigParse(format!("{context} is not UTF-8 JSON: {e}")))
1876}
1877
1878#[cfg(test)]
1883mod tests {
1884 use tempfile::tempdir;
1885
1886 use oci_client::manifest::{ImageIndexEntry, Platform as OciPlatform};
1887
1888 use super::{
1889 LayerDescriptor, MaterializeLayersRequest, Platform, layer_work_path,
1890 read_erofs_layer_metadata, resolve_cached_metadata_pull_result_async,
1891 resolve_cached_pull_result, resolve_platform_digest,
1892 };
1893 use crate::{
1894 cache::{CachedImageMetadata, CachedLayerMetadata, GlobalCache},
1895 config::ImageConfig,
1896 digest::Digest,
1897 erofs,
1898 error::ImageError,
1899 ext4::EXT4_ROOTFS_MATERIALIZER_ABI,
1900 pull::{PullOptions, PullPolicy, RootfsMaterialization},
1901 tree::{
1902 FileData, FileTree, InodeMetadata, RegularFileId, RegularFileNode, TreeNode, Xattr,
1903 },
1904 };
1905
1906 #[test]
1907 fn test_platform_resolver_prefers_exact_variant() {
1908 let manifests = vec![
1909 ImageIndexEntry {
1910 media_type: "application/vnd.oci.image.manifest.v1+json".into(),
1911 digest: "sha256:arch-only".into(),
1912 size: 1,
1913 platform: Some(OciPlatform {
1914 architecture: "arm".into(),
1915 os: "linux".into(),
1916 os_version: None,
1917 os_features: None,
1918 variant: None,
1919 features: None,
1920 }),
1921 annotations: None,
1922 artifact_type: None,
1923 },
1924 ImageIndexEntry {
1925 media_type: "application/vnd.oci.image.manifest.v1+json".into(),
1926 digest: "sha256:exact".into(),
1927 size: 1,
1928 platform: Some(OciPlatform {
1929 architecture: "arm".into(),
1930 os: "linux".into(),
1931 os_version: None,
1932 os_features: None,
1933 variant: Some("v7".into()),
1934 features: None,
1935 }),
1936 annotations: None,
1937 artifact_type: None,
1938 },
1939 ];
1940
1941 let digest =
1942 resolve_platform_digest(&manifests, &Platform::with_variant("linux", "arm", "v7"));
1943 assert_eq!(digest.as_deref(), Some("sha256:exact"));
1944 }
1945
1946 #[test]
1947 fn test_layer_work_path_uses_path_safe_digest() {
1948 let temp = tempdir().unwrap();
1949 let digest = Digest::new("sha256", "abc123");
1950
1951 let path = layer_work_path(temp.path(), &digest, "erofs.part");
1952 let file_name = path.file_name().unwrap().to_string_lossy();
1953
1954 assert_eq!(file_name, "sha256_abc123.erofs.part");
1955 assert!(!file_name.contains(':'));
1956 }
1957
1958 #[test]
1959 fn test_resolve_cached_pull_result_if_missing_uses_complete_cache() {
1960 let temp = tempdir().unwrap();
1961 let cache = GlobalCache::new(temp.path()).unwrap();
1962 let reference: oci_client::Reference = "docker.io/library/alpine".parse().unwrap();
1963 let metadata = write_cached_image_fixture(&cache, &reference, &[true, true]);
1964
1965 let cached = resolve_cached_pull_result(
1966 &cache,
1967 &reference,
1968 &PullOptions {
1969 pull_policy: PullPolicy::IfMissing,
1970 force: false,
1971 ..Default::default()
1972 },
1973 )
1974 .unwrap()
1975 .expect("expected cached pull result");
1976
1977 assert!(cached.result.cached);
1978 assert_eq!(cached.result.layer_diff_ids.len(), 2);
1979 assert_eq!(
1980 cached.result.manifest_digest.to_string(),
1981 metadata.manifest_digest
1982 );
1983 assert_eq!(cached.result.config.env, metadata.config.env);
1984 assert_eq!(
1985 cached.result.layer_diff_ids[0].to_string(),
1986 metadata.layers[0].diff_id
1987 );
1988 assert_eq!(
1989 cached.result.layer_diff_ids[1].to_string(),
1990 metadata.layers[1].diff_id
1991 );
1992 }
1993
1994 #[test]
1995 fn test_resolve_cached_pull_result_never_uses_complete_cache() {
1996 let temp = tempdir().unwrap();
1997 let cache = GlobalCache::new(temp.path()).unwrap();
1998 let reference: oci_client::Reference = "docker.io/library/busybox:latest".parse().unwrap();
1999 write_cached_image_fixture(&cache, &reference, &[true]);
2000
2001 let cached = resolve_cached_pull_result(
2002 &cache,
2003 &reference,
2004 &PullOptions {
2005 pull_policy: PullPolicy::Never,
2006 force: false,
2007 ..Default::default()
2008 },
2009 )
2010 .unwrap();
2011
2012 assert!(cached.is_some());
2013 assert!(cached.unwrap().result.cached);
2014 }
2015
2016 #[test]
2017 fn test_pull_cached_uses_complete_cache() {
2018 let temp = tempdir().unwrap();
2019 let cache = GlobalCache::new(temp.path()).unwrap();
2020 let reference: oci_client::Reference = "docker.io/library/alpine".parse().unwrap();
2021 let metadata = write_cached_image_fixture(&cache, &reference, &[true]);
2022
2023 let cached = super::Registry::pull_cached(
2024 &cache,
2025 &reference,
2026 &PullOptions {
2027 pull_policy: PullPolicy::IfMissing,
2028 force: false,
2029 ..Default::default()
2030 },
2031 )
2032 .unwrap()
2033 .expect("expected cached pull result");
2034
2035 assert!(cached.0.cached);
2036 assert_eq!(
2037 cached.0.manifest_digest.to_string(),
2038 metadata.manifest_digest
2039 );
2040 assert_eq!(cached.1.manifest_digest, metadata.manifest_digest);
2041 }
2042
2043 #[tokio::test]
2044 async fn test_pull_cached_by_manifest_digest_finds_tag_keyed_metadata() {
2045 let temp = tempdir().unwrap();
2046 let cache = GlobalCache::new(temp.path()).unwrap();
2047 let reference: oci_client::Reference = "docker.io/library/python:3.12".parse().unwrap();
2048 let metadata = write_cached_image_fixture(&cache, &reference, &[true, true]);
2049 let manifest_digest = parse_digest(&metadata.manifest_digest);
2050
2051 let cached = super::Registry::pull_cached_by_manifest_digest(&cache, &manifest_digest)
2052 .await
2053 .unwrap()
2054 .expect("expected cached pull result by manifest digest");
2055
2056 assert!(cached.0.cached);
2057 assert_eq!(
2058 cached.0.manifest_digest.to_string(),
2059 metadata.manifest_digest
2060 );
2061 assert_eq!(cached.1.manifest_digest, metadata.manifest_digest);
2062 }
2063
2064 #[tokio::test]
2065 async fn snapshot_flat_metadata_does_not_require_disks_or_accept_moved_tags() {
2066 let temp = tempdir().unwrap();
2067 let cache = GlobalCache::new(temp.path()).unwrap();
2068 let reference: oci_client::Reference = "docker.io/library/alpine:latest".parse().unwrap();
2069 let metadata = write_cached_image_fixture(&cache, &reference, &[false, false]);
2070 let digest = parse_digest(&metadata.manifest_digest);
2071 for refs in [vec![reference.clone()], vec![]] {
2072 let found = super::Registry::pull_snapshot_cached(
2073 &cache,
2074 &refs,
2075 &digest,
2076 RootfsMaterialization::Flat,
2077 )
2078 .await
2079 .unwrap()
2080 .unwrap();
2081 assert_eq!(found.0.manifest_digest, digest);
2082 assert_eq!(found.0.config.env, metadata.config.env);
2083 assert!(
2084 super::Registry::pull_snapshot_cached(
2085 &cache,
2086 &refs,
2087 &digest,
2088 RootfsMaterialization::Layered
2089 )
2090 .await
2091 .unwrap()
2092 .is_none()
2093 );
2094 }
2095 let other = parse_digest(&format!("sha256:{}", "f".repeat(64)));
2096 assert!(
2097 super::Registry::pull_snapshot_cached(
2098 &cache,
2099 &[reference],
2100 &other,
2101 RootfsMaterialization::Flat
2102 )
2103 .await
2104 .unwrap()
2105 .is_none()
2106 );
2107 }
2108
2109 #[tokio::test]
2110 async fn test_pull_cached_by_manifest_digest_requires_complete_artifacts() {
2111 let temp = tempdir().unwrap();
2112 let cache = GlobalCache::new(temp.path()).unwrap();
2113 let reference: oci_client::Reference = "docker.io/library/python:3.12".parse().unwrap();
2114 let metadata = write_cached_image_fixture(&cache, &reference, &[true, false]);
2115 let manifest_digest = parse_digest(&metadata.manifest_digest);
2116
2117 let cached = super::Registry::pull_cached_by_manifest_digest(&cache, &manifest_digest)
2118 .await
2119 .unwrap();
2120
2121 assert!(cached.is_none());
2122 }
2123
2124 #[tokio::test]
2125 async fn test_pull_never_returns_not_cached_when_any_layer_is_missing() {
2126 let temp = tempdir().unwrap();
2127 let cache = GlobalCache::new(temp.path()).unwrap();
2128 let reference: oci_client::Reference = "docker.io/library/debian:stable".parse().unwrap();
2129 write_cached_image_fixture(&cache, &reference, &[true, false]);
2130
2131 let cached = resolve_cached_pull_result(
2132 &cache,
2133 &reference,
2134 &PullOptions {
2135 pull_policy: PullPolicy::Never,
2136 force: false,
2137 ..Default::default()
2138 },
2139 )
2140 .unwrap();
2141 assert!(cached.is_none());
2142
2143 let registry = super::Registry::new(Platform::default(), cache).unwrap();
2144 let err = registry
2145 .pull(
2146 &reference,
2147 &PullOptions {
2148 pull_policy: PullPolicy::Never,
2149 force: false,
2150 ..Default::default()
2151 },
2152 )
2153 .await;
2154
2155 assert!(matches!(err, Err(ImageError::NotCached { .. })));
2156 }
2157
2158 #[test]
2159 fn test_resolve_cached_pull_result_ignores_corrupt_metadata_file() {
2160 let temp = tempdir().unwrap();
2161 let cache = GlobalCache::new(temp.path()).unwrap();
2162 let reference: oci_client::Reference = "docker.io/library/ubuntu:latest".parse().unwrap();
2163 let metadata_path = image_metadata_path(temp.path(), &reference);
2164 std::fs::write(&metadata_path, b"{ definitely not json").unwrap();
2165
2166 let cached = resolve_cached_pull_result(
2167 &cache,
2168 &reference,
2169 &PullOptions {
2170 pull_policy: PullPolicy::IfMissing,
2171 force: false,
2172 ..Default::default()
2173 },
2174 )
2175 .unwrap();
2176
2177 assert!(cached.is_none());
2178 }
2179
2180 #[test]
2181 fn test_resolve_cached_pull_result_skips_cache_for_force_and_always() {
2182 let temp = tempdir().unwrap();
2183 let cache = GlobalCache::new(temp.path()).unwrap();
2184 let reference: oci_client::Reference = "docker.io/library/fedora:latest".parse().unwrap();
2185 write_cached_image_fixture(&cache, &reference, &[true]);
2186
2187 let forced = resolve_cached_pull_result(
2188 &cache,
2189 &reference,
2190 &PullOptions {
2191 pull_policy: PullPolicy::IfMissing,
2192 force: true,
2193 ..Default::default()
2194 },
2195 )
2196 .unwrap();
2197 assert!(forced.is_none());
2198
2199 let always = resolve_cached_pull_result(
2200 &cache,
2201 &reference,
2202 &PullOptions {
2203 pull_policy: PullPolicy::Always,
2204 force: false,
2205 ..Default::default()
2206 },
2207 )
2208 .unwrap();
2209 assert!(always.is_none());
2210 }
2211
2212 #[test]
2213 fn test_resolve_cached_pull_result_ignores_invalid_digest_metadata() {
2214 let temp = tempdir().unwrap();
2215 let cache = GlobalCache::new(temp.path()).unwrap();
2216 let reference: oci_client::Reference = "docker.io/library/redis:latest".parse().unwrap();
2217 let mut metadata = write_cached_image_fixture(&cache, &reference, &[true]);
2218 metadata.layers[0].diff_id = "not-a-digest".into();
2219 std::fs::write(
2221 cache.image_metadata_path(&reference),
2222 serde_json::to_vec(&metadata).unwrap(),
2223 )
2224 .unwrap();
2225
2226 let cached = resolve_cached_pull_result(
2227 &cache,
2228 &reference,
2229 &PullOptions {
2230 pull_policy: PullPolicy::IfMissing,
2231 force: false,
2232 ..Default::default()
2233 },
2234 )
2235 .unwrap();
2236
2237 assert!(cached.is_none());
2238 }
2239
2240 #[test]
2241 fn test_resolve_cached_pull_result_requires_fsmeta_and_vmdk() {
2242 let temp = tempdir().unwrap();
2243 let cache = GlobalCache::new(temp.path()).unwrap();
2244 let reference: oci_client::Reference = "docker.io/library/alpine:latest".parse().unwrap();
2245 let metadata = write_cached_image_fixture(&cache, &reference, &[false, false]);
2247 let manifest_digest = parse_digest(&metadata.manifest_digest);
2248 for (index, _) in metadata.layers.iter().enumerate() {
2250 let diff_id = parse_digest(&format!("sha256:{:064x}", index as u64 + 1000));
2251 write_valid_erofs_layer(&cache, &diff_id, b"cached fixture layer");
2252 }
2253 let _ = std::fs::remove_file(cache.fsmeta_erofs_path(&manifest_digest));
2255 let _ = std::fs::remove_file(cache.vmdk_path(&manifest_digest));
2256
2257 let cached = resolve_cached_pull_result(
2258 &cache,
2259 &reference,
2260 &PullOptions {
2261 pull_policy: PullPolicy::IfMissing,
2262 force: false,
2263 ..Default::default()
2264 },
2265 )
2266 .unwrap();
2267
2268 assert!(cached.is_none(), "should not be cached without fsmeta+VMDK");
2269 }
2270
2271 #[tokio::test]
2272 async fn test_pull_never_treats_invalid_digest_metadata_as_not_cached() {
2273 let temp = tempdir().unwrap();
2274 let cache = GlobalCache::new(temp.path()).unwrap();
2275 let reference: oci_client::Reference = "docker.io/library/httpd:latest".parse().unwrap();
2276 let mut metadata = write_cached_image_fixture(&cache, &reference, &[true]);
2277 metadata.layers[0].diff_id = "not-a-digest".into();
2278 std::fs::write(
2280 cache.image_metadata_path(&reference),
2281 serde_json::to_vec(&metadata).unwrap(),
2282 )
2283 .unwrap();
2284
2285 let registry = super::Registry::new(Platform::default(), cache).unwrap();
2286 let result = registry
2287 .pull(
2288 &reference,
2289 &PullOptions {
2290 pull_policy: PullPolicy::Never,
2291 force: false,
2292 ..Default::default()
2293 },
2294 )
2295 .await;
2296
2297 assert!(matches!(result, Err(ImageError::NotCached { .. })));
2298 }
2299
2300 #[tokio::test]
2301 async fn test_pull_with_progress_cached_if_missing_emits_only_summary_events() {
2302 let temp = tempdir().unwrap();
2303 let cache = GlobalCache::new(temp.path()).unwrap();
2304 let reference: oci_client::Reference = "docker.io/library/nginx:latest".parse().unwrap();
2305 write_cached_image_fixture(&cache, &reference, &[true, true]);
2306 let registry = super::Registry::new(Platform::default(), cache).unwrap();
2307
2308 let (mut handle, task) = registry.pull_with_progress(
2309 &reference,
2310 &PullOptions {
2311 pull_policy: PullPolicy::IfMissing,
2312 force: false,
2313 ..Default::default()
2314 },
2315 );
2316
2317 let result = task.await.unwrap().unwrap();
2318 let mut events = Vec::new();
2319 while let Some(event) = handle.recv().await {
2320 events.push(event);
2321 }
2322
2323 assert!(result.cached);
2324 assert_eq!(events.len(), 3);
2325 assert!(matches!(
2326 &events[0],
2327 crate::progress::PullProgress::Resolving { reference: event_ref }
2328 if event_ref.as_ref() == reference.to_string()
2329 ));
2330 assert!(matches!(
2331 &events[1],
2332 crate::progress::PullProgress::Resolved {
2333 reference: event_ref,
2334 layer_count: 2,
2335 ..
2336 } if event_ref.as_ref() == reference.to_string()
2337 ));
2338 assert!(matches!(
2339 &events[2],
2340 crate::progress::PullProgress::Complete {
2341 reference: event_ref,
2342 layer_count: 2,
2343 } if event_ref.as_ref() == reference.to_string()
2344 ));
2345 }
2346
2347 #[tokio::test]
2348 async fn test_flat_cache_hit_does_not_require_layered_outputs() {
2349 let temp = tempdir().unwrap();
2350 let cache = GlobalCache::new(temp.path()).unwrap();
2351 let reference: oci_client::Reference =
2352 "docker.io/library/shared-base:flat".parse().unwrap();
2353 let metadata = write_cached_image_fixture(&cache, &reference, &[true]);
2354 let diff_id = parse_digest(&metadata.layers[0].diff_id);
2355 write_valid_erofs_layer(&cache, &diff_id, b"shared flat contents");
2356 let manifest_digest = parse_digest(&metadata.manifest_digest);
2357 let registry = super::Registry::new(Platform::default(), cache.clone()).unwrap();
2358 registry
2359 .materialize_flat_rootfs(&manifest_digest, std::slice::from_ref(&diff_id), false)
2360 .await
2361 .unwrap();
2362 std::fs::remove_file(cache.fsmeta_erofs_path(&manifest_digest)).unwrap();
2363 std::fs::remove_file(cache.vmdk_path(&manifest_digest)).unwrap();
2364
2365 let flat = resolve_cached_pull_result(
2366 &cache,
2367 &reference,
2368 &PullOptions {
2369 materialization: RootfsMaterialization::Flat,
2370 ..Default::default()
2371 },
2372 )
2373 .unwrap();
2374 let layered = resolve_cached_pull_result(
2375 &cache,
2376 &reference,
2377 &PullOptions {
2378 materialization: RootfsMaterialization::Layered,
2379 ..Default::default()
2380 },
2381 )
2382 .unwrap();
2383 let all = resolve_cached_pull_result(
2384 &cache,
2385 &reference,
2386 &PullOptions {
2387 materialization: RootfsMaterialization::All,
2388 ..Default::default()
2389 },
2390 )
2391 .unwrap();
2392
2393 assert!(flat.is_some(), "flat should not depend on fsmeta or VMDK");
2394 assert!(layered.is_none(), "layered still requires fsmeta and VMDK");
2395 assert!(all.is_none(), "all requires both representations");
2396 }
2397
2398 #[tokio::test]
2399 async fn test_flat_cache_hit_rejects_stale_materializer_inputs() {
2400 let temp = tempdir().unwrap();
2401 let cache = GlobalCache::new(temp.path()).unwrap();
2402 let reference: oci_client::Reference =
2403 "docker.io/library/shared-base:stale-flat".parse().unwrap();
2404 let metadata = write_cached_image_fixture(&cache, &reference, &[true]);
2405 let diff_id = parse_digest(&metadata.layers[0].diff_id);
2406 write_valid_erofs_layer(&cache, &diff_id, b"shared flat contents");
2407 let manifest_digest = parse_digest(&metadata.manifest_digest);
2408 let platform = Platform::default();
2409 let registry = super::Registry::new(platform.clone(), cache.clone()).unwrap();
2410 let current = registry
2411 .materialize_flat_rootfs(&manifest_digest, std::slice::from_ref(&diff_id), false)
2412 .await
2413 .unwrap();
2414 let options = PullOptions {
2415 materialization: RootfsMaterialization::Flat,
2416 ..Default::default()
2417 };
2418
2419 let mut stale_abi = current.clone();
2420 stale_abi.materializer_abi = EXT4_ROOTFS_MATERIALIZER_ABI.saturating_sub(1);
2421 cache.write_flat_ref(&manifest_digest, &stale_abi).unwrap();
2422 assert!(
2423 resolve_cached_pull_result(&cache, &reference, &options)
2424 .unwrap()
2425 .is_none(),
2426 "the synchronous cache path must reject an older materializer ABI"
2427 );
2428 assert!(
2429 resolve_cached_metadata_pull_result_async(
2430 &cache,
2431 metadata.clone(),
2432 RootfsMaterialization::Flat,
2433 &platform,
2434 )
2435 .await
2436 .unwrap()
2437 .is_none(),
2438 "the asynchronous cache path must reject an older materializer ABI"
2439 );
2440
2441 let mut stale_derivation = current;
2442 stale_derivation.derivation_digest = format!("sha256:{}", "f".repeat(64));
2443 cache
2444 .write_flat_ref(&manifest_digest, &stale_derivation)
2445 .unwrap();
2446 assert!(
2447 resolve_cached_pull_result(&cache, &reference, &options)
2448 .unwrap()
2449 .is_none(),
2450 "the synchronous cache path must reject a different derivation"
2451 );
2452 assert!(
2453 resolve_cached_metadata_pull_result_async(
2454 &cache,
2455 metadata,
2456 RootfsMaterialization::Flat,
2457 &platform,
2458 )
2459 .await
2460 .unwrap()
2461 .is_none(),
2462 "the asynchronous cache path must reject a different derivation"
2463 );
2464 }
2465
2466 #[tokio::test]
2467 async fn test_flat_layer_stage_reuses_erofs_and_skips_layered_outputs() {
2468 let temp = tempdir().unwrap();
2469 let cache = GlobalCache::new(temp.path()).unwrap();
2470 let registry = super::Registry::new(Platform::default(), cache.clone()).unwrap();
2471 let reference: oci_client::Reference =
2472 "docker.io/library/shared-base:flat".parse().unwrap();
2473 let manifest_digest = parse_digest(&format!("sha256:{}", "a".repeat(64)));
2474 let diff_id = parse_digest(&format!("sha256:{}", "b".repeat(64)));
2475 let descriptor_digest = parse_digest(&format!("sha256:{}", "c".repeat(64)));
2476 let tar_path = cache.tar_path(&descriptor_digest);
2477 write_valid_erofs_layer(&cache, &diff_id, b"cached shared layer");
2478
2479 registry
2480 .materialize_layer_stage(MaterializeLayersRequest {
2481 oci_ref: &reference,
2482 manifest_digest: &manifest_digest,
2483 layer_descriptors: &[LayerDescriptor {
2484 digest: descriptor_digest,
2485 media_type: Some("application/vnd.oci.image.layer.v1.tar+gzip".into()),
2486 size: Some(1024),
2487 }],
2488 diff_ids: &[diff_id.to_string()],
2489 force: false,
2490 materialization: RootfsMaterialization::Flat,
2491 progress: None,
2492 staged_layers: None,
2493 })
2494 .await
2495 .unwrap();
2496
2497 assert!(!cache.is_fsmeta_materialized(&manifest_digest));
2498 assert!(!cache.is_vmdk_materialized(&manifest_digest));
2499 assert!(!tar_path.exists());
2500 }
2501
2502 #[tokio::test]
2503 async fn test_layered_target_rebuilds_fsmeta_from_cached_erofs() {
2504 let temp = tempdir().unwrap();
2505 let cache = GlobalCache::new(temp.path()).unwrap();
2506 let registry = super::Registry::new(Platform::default(), cache.clone()).unwrap();
2507 let reference: oci_client::Reference =
2508 "docker.io/library/shared-base:layered".parse().unwrap();
2509 let manifest_digest = parse_digest(&format!("sha256:{}", "d".repeat(64)));
2510 let diff_id = parse_digest(&format!("sha256:{}", "e".repeat(64)));
2511 let descriptor_digest = parse_digest(&format!("sha256:{}", "f".repeat(64)));
2512 let tar_path = cache.tar_path(&descriptor_digest);
2513 write_valid_erofs_layer(&cache, &diff_id, b"reused without source tar");
2514
2515 registry
2516 .materialize_layer_stage(MaterializeLayersRequest {
2517 oci_ref: &reference,
2518 manifest_digest: &manifest_digest,
2519 layer_descriptors: &[LayerDescriptor {
2520 digest: descriptor_digest,
2521 media_type: Some("application/vnd.oci.image.layer.v1.tar+gzip".into()),
2522 size: Some(2048),
2523 }],
2524 diff_ids: &[diff_id.to_string()],
2525 force: false,
2526 materialization: RootfsMaterialization::Layered,
2527 progress: None,
2528 staged_layers: None,
2529 })
2530 .await
2531 .unwrap();
2532
2533 assert!(cache.is_fsmeta_materialized(&manifest_digest));
2534 assert!(cache.is_vmdk_materialized(&manifest_digest));
2535 assert!(!tar_path.exists());
2536
2537 let mut reader = erofs::ErofsReader::new(
2538 std::fs::File::open(cache.fsmeta_erofs_path(&manifest_digest)).unwrap(),
2539 )
2540 .unwrap();
2541 let file = reader.inode_debug_info("/shared.txt").unwrap();
2542 assert_eq!(file.size, b"reused without source tar".len() as u64);
2543 }
2544
2545 #[test]
2546 fn test_erofs_metadata_reconstruction_preserves_hardlinks_and_blocks() {
2547 let temp = tempdir().unwrap();
2548 let path = temp.path().join("layer.erofs");
2549 let id = RegularFileId::new();
2550 let node = || {
2551 TreeNode::RegularFile(RegularFileNode {
2552 id,
2553 metadata: InodeMetadata::default(),
2554 xattrs: Vec::new(),
2555 data: FileData::Memory(b"shared hardlink data".to_vec()),
2556 nlink: 2,
2557 })
2558 };
2559 let mut tree = FileTree::new();
2560 tree.root.metadata.uid = 42;
2561 tree.root.xattrs.push(Xattr {
2562 name: b"user.root-test".to_vec(),
2563 value: b"preserved".to_vec(),
2564 });
2565 tree.insert(b"nested/alpha", node()).unwrap();
2566 tree.insert(b"nested/beta", node()).unwrap();
2567 erofs::write_erofs(&tree, &path).unwrap();
2568
2569 let (reconstructed, data_map) = read_erofs_layer_metadata(&path).unwrap();
2570 let TreeNode::RegularFile(alpha) = reconstructed.get(b"nested/alpha").unwrap() else {
2571 panic!("alpha should be a regular file");
2572 };
2573 let TreeNode::RegularFile(beta) = reconstructed.get(b"nested/beta").unwrap() else {
2574 panic!("beta should be a regular file");
2575 };
2576
2577 assert_eq!(alpha.id, beta.id);
2578 assert_eq!(reconstructed.root.metadata.uid, 42);
2579 assert_eq!(reconstructed.root.xattrs.len(), 1);
2580 assert_eq!(reconstructed.root.xattrs[0].name, b"user.root-test");
2581 assert_eq!(reconstructed.root.xattrs[0].value, b"preserved");
2582 assert_eq!(
2583 data_map
2584 .file_blocks
2585 .get(std::path::Path::new("nested/alpha")),
2586 data_map
2587 .file_blocks
2588 .get(std::path::Path::new("nested/beta"))
2589 );
2590 assert_eq!(
2591 data_map.file_blocks[std::path::Path::new("nested/alpha")].1,
2592 b"shared hardlink data".len() as u64
2593 );
2594 }
2595
2596 fn write_valid_erofs_layer(cache: &GlobalCache, diff_id: &Digest, contents: &[u8]) {
2597 let mut tree = FileTree::new();
2598 tree.insert(
2599 b"shared.txt",
2600 TreeNode::RegularFile(RegularFileNode {
2601 id: RegularFileId::new(),
2602 metadata: InodeMetadata::default(),
2603 xattrs: Vec::new(),
2604 data: FileData::Memory(contents.to_vec()),
2605 nlink: 1,
2606 }),
2607 )
2608 .unwrap();
2609 erofs::write_erofs(&tree, &cache.layer_erofs_path(diff_id)).unwrap();
2610 }
2611
2612 fn write_cached_image_fixture(
2613 cache: &GlobalCache,
2614 reference: &oci_client::Reference,
2615 materialized_layers: &[bool],
2616 ) -> CachedImageMetadata {
2617 let metadata = CachedImageMetadata {
2618 manifest_digest:
2619 "sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
2620 .to_string(),
2621 config_digest:
2622 "sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
2623 .to_string(),
2624 raw_manifest_json: r#"{"schemaVersion":2,"layers":[]}"#.to_string(),
2625 raw_config_json:
2626 r#"{"architecture":"amd64","os":"linux","rootfs":{"type":"layers","diff_ids":[]}}"#
2627 .to_string(),
2628 config: ImageConfig {
2629 env: vec!["PATH=/usr/bin".into()],
2630 ..Default::default()
2631 },
2632 layers: materialized_layers
2633 .iter()
2634 .enumerate()
2635 .map(|(index, _)| CachedLayerMetadata {
2636 digest: layer_digest(index),
2637 media_type: Some("application/vnd.oci.image.layer.v1.tar+gzip".into()),
2638 size_bytes: Some((index as u64 + 1) * 100),
2639 diff_id: format!("sha256:{:064x}", index as u64 + 1000),
2640 })
2641 .collect(),
2642 };
2643
2644 cache.write_image_metadata(reference, &metadata).unwrap();
2645
2646 let all_materialized = materialized_layers.iter().all(|m| *m);
2648 for (index, materialized) in materialized_layers.iter().copied().enumerate() {
2649 let diff_id = parse_digest(&format!("sha256:{:064x}", index as u64 + 1000));
2650 if materialized {
2651 write_valid_erofs_layer(cache, &diff_id, b"cached fixture layer");
2652 }
2653 }
2654
2655 if all_materialized && !materialized_layers.is_empty() {
2657 let manifest_digest = parse_digest(&metadata.manifest_digest);
2658 erofs::write_erofs(&FileTree::new(), &cache.fsmeta_erofs_path(&manifest_digest))
2659 .unwrap();
2660 std::fs::write(cache.vmdk_path(&manifest_digest), b"# VMDK fixture").unwrap();
2661 }
2662
2663 metadata
2664 }
2665
2666 fn layer_digest(index: usize) -> String {
2667 format!("sha256:{:064x}", index as u64 + 1)
2668 }
2669
2670 fn parse_digest(digest: &str) -> crate::digest::Digest {
2671 digest.parse().unwrap()
2672 }
2673
2674 fn image_metadata_path(
2675 cache_root: &std::path::Path,
2676 reference: &oci_client::Reference,
2677 ) -> std::path::PathBuf {
2678 use sha2::{Digest as Sha2Digest, Sha256};
2679
2680 let mut hasher = Sha256::new();
2681 hasher.update(reference.to_string().as_bytes());
2682 cache_root
2683 .join("manifests")
2684 .join(format!("{}.json", hex::encode(hasher.finalize())))
2685 }
2686
2687 #[test]
2688 fn test_registry_builder_default() {
2689 let temp = tempdir().unwrap();
2690 let cache = GlobalCache::new(temp.path()).unwrap();
2691 let registry = super::Registry::builder(Platform::default(), cache)
2692 .build()
2693 .unwrap();
2694
2695 assert!(matches!(
2696 registry.auth,
2697 oci_client::secrets::RegistryAuth::Anonymous
2698 ));
2699 }
2700
2701 #[test]
2702 fn test_registry_builder_with_auth() {
2703 let temp = tempdir().unwrap();
2704 let cache = GlobalCache::new(temp.path()).unwrap();
2705 let registry = super::Registry::builder(Platform::default(), cache)
2706 .auth(crate::RegistryAuth::Basic {
2707 username: "user".into(),
2708 password: "pass".into(),
2709 })
2710 .build()
2711 .unwrap();
2712
2713 assert!(matches!(
2714 registry.auth,
2715 oci_client::secrets::RegistryAuth::Basic(_, _)
2716 ));
2717 }
2718
2719 #[test]
2720 fn test_registry_builder_with_insecure_registries() {
2721 let temp = tempdir().unwrap();
2722 let cache = GlobalCache::new(temp.path()).unwrap();
2723 super::Registry::builder(Platform::default(), cache)
2726 .add_insecure_registries(vec!["localhost:5000".into()])
2727 .build()
2728 .unwrap();
2729 }
2730
2731 fn generate_test_ca_pem() -> Vec<u8> {
2733 let key_pair = rcgen::KeyPair::generate().unwrap();
2734 let mut params = rcgen::CertificateParams::default();
2735 params.is_ca = rcgen::IsCa::Ca(rcgen::BasicConstraints::Unconstrained);
2736 let cert = params.self_signed(&key_pair).unwrap();
2737 cert.pem().into_bytes()
2738 }
2739
2740 #[test]
2741 fn test_registry_builder_with_valid_ca_cert() {
2742 let temp = tempdir().unwrap();
2743 let cache = GlobalCache::new(temp.path()).unwrap();
2744 let pem = generate_test_ca_pem();
2745 super::Registry::builder(Platform::default(), cache)
2746 .extra_ca_certs(vec![pem])
2747 .build()
2748 .unwrap();
2749 }
2750
2751 fn build_err(result: Result<super::Registry, crate::ImageError>) -> crate::ImageError {
2753 match result {
2754 Err(e) => e,
2755 Ok(_) => panic!("expected build to fail"),
2756 }
2757 }
2758
2759 #[test]
2760 fn test_registry_builder_rejects_invalid_pem() {
2761 let temp = tempdir().unwrap();
2762 let cache = GlobalCache::new(temp.path()).unwrap();
2763 let bad_pem = b"not valid PEM data".to_vec();
2764 let err = build_err(
2765 super::Registry::builder(Platform::default(), cache)
2766 .extra_ca_certs(vec![bad_pem])
2767 .build(),
2768 );
2769
2770 assert!(
2771 err.to_string().contains("no certificates found"),
2772 "expected 'no certificates found', got: {err}"
2773 );
2774 }
2775
2776 #[test]
2777 fn test_registry_builder_rejects_empty_pem() {
2778 let temp = tempdir().unwrap();
2779 let cache = GlobalCache::new(temp.path()).unwrap();
2780 let err = build_err(
2781 super::Registry::builder(Platform::default(), cache)
2782 .extra_ca_certs(vec![Vec::new()])
2783 .build(),
2784 );
2785
2786 assert!(
2787 err.to_string().contains("no certificates found"),
2788 "expected 'no certificates found', got: {err}"
2789 );
2790 }
2791
2792 #[test]
2793 fn test_registry_builder_all_options() {
2794 let temp = tempdir().unwrap();
2795 let cache = GlobalCache::new(temp.path()).unwrap();
2796 let pem = generate_test_ca_pem();
2797 super::Registry::builder(Platform::default(), cache)
2798 .auth(crate::RegistryAuth::Basic {
2799 username: "user".into(),
2800 password: "pass".into(),
2801 })
2802 .add_insecure_registries(vec!["localhost:5000".into()])
2803 .extra_ca_certs(vec![pem])
2804 .build()
2805 .unwrap();
2806 }
2807
2808 #[test]
2809 fn test_registry_new_equals_builder_default() {
2810 let temp = tempdir().unwrap();
2811 let cache = GlobalCache::new(temp.path()).unwrap();
2812 super::Registry::new(Platform::default(), cache).unwrap();
2814 }
2815}