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