Skip to main content

microsandbox_image/registry/
client.rs

1//! OCI registry client.
2//!
3//! Wraps `oci-client` with platform resolution, caching, and progress reporting.
4
5use 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
44//--------------------------------------------------------------------------------------------------
45// Constants
46//--------------------------------------------------------------------------------------------------
47
48/// Minimum byte delta between per-layer materialization progress updates.
49const MATERIALIZE_PROGRESS_EMIT_BYTES: u64 = 256 * 1024;
50
51/// Upper bound for concurrently active layer download/materialize tasks.
52const MAX_LAYER_PIPELINE_CONCURRENCY: usize = 16;
53
54//--------------------------------------------------------------------------------------------------
55// Types
56//--------------------------------------------------------------------------------------------------
57
58/// OCI registry client with platform resolution, caching, and progress reporting.
59pub 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/// Resolved manifest layer descriptor used during download/materialization.
67#[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
94/// Per-layer pipeline success: EROFS image written, data-stripped tree + data map retained.
95/// `tree` and `data_map` are `None` when the EROFS was already cached.
96struct LayerPipelineTreeSuccess {
97    layer_index: usize,
98    tree: Option<FileTree>,
99    data_map: Option<erofs::ErofsDataMap>,
100}
101
102/// Wraps an `AsyncRead` to emit `LayerMaterializeProgress` events as the
103/// tar stream is read during EROFS materialization.
104///
105/// Progress events are throttled to avoid flooding the channel — an update
106/// is sent only after at least `MATERIALIZE_PROGRESS_EMIT_BYTES` (256 KiB)
107/// have been read since the last event.
108struct 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
117//--------------------------------------------------------------------------------------------------
118// Methods
119//--------------------------------------------------------------------------------------------------
120
121impl<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    /// Create a registry client with anonymous authentication and default TLS settings.
141    pub fn new(platform: Platform, cache: GlobalCache) -> ImageResult<Self> {
142        Self::builder(platform, cache).build()
143    }
144
145    /// Create a builder for configuring auth, TLS, and other registry options.
146    pub fn builder(platform: Platform, cache: GlobalCache) -> RegistryBuilder {
147        RegistryBuilder::new(platform, cache)
148    }
149
150    /// Resolve a pull directly from the on-disk cache without building a registry client.
151    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    /// Resolve a pull from any complete cached image metadata matching a manifest digest.
161    ///
162    /// Snapshot restores use this before contacting a registry because the
163    /// snapshot already records the immutable base image digest, while cached
164    /// metadata may still be keyed by the mutable tag that produced it.
165    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    /// Resolve snapshot image defaults by immutable digest. Flat snapshots own
177    /// their complete disk, so they need metadata but no materialized OCI layers.
178    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        // Most snapshots retain either their original tag key or a pinned key.
187        // Avoid scanning every unrelated image on this common path.
188        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    /// Fetch only the immutable manifest and config needed by a flat snapshot.
207    /// Normal pulls still independently require their filesystem artifacts.
208    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    /// Pull an image. Downloads blobs and materializes EROFS layers concurrently.
258    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    /// Materialize layers that have already been staged into the cache.
267    ///
268    /// This is used by archive import: layer blobs are written to the cache at
269    /// their descriptor digest first, then the normal EROFS/fsmeta/VMDK
270    /// materialization path can run without contacting a registry.
271    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    /// Materialize a reusable flat ext4 rootfs from already cached EROFS layers.
282    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    /// Pull with progress reporting.
375    ///
376    /// Creates a progress channel internally and returns both the receiver
377    /// handle and the spawned pull task.
378    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    /// Pull with an externally-provided progress sender.
392    ///
393    /// Use [`progress_channel()`](crate::progress_channel) to create the
394    /// channel, keep the [`PullProgressHandle`] receiver, and pass the
395    /// [`PullProgressSender`] here.
396    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    /// Spawn the pull task with a progress sender.
409    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    /// Core pull implementation.
442    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        // Lock files are intentionally never deleted — stable inodes prevent
460        // TOCTOU races where two processes flock different inodes at the same path.
461        let _image_lock_guard = scopeguard::guard(image_lock_file, |file| {
462            let _ = flock_unlock(&file);
463        });
464
465        // Step 1: Early cache check using persisted image metadata.
466        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        // Step 2: Resolve manifest.
507        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        // Determine media type from manifest bytes. For multi-platform images,
520        // this also fetches the platform-specific config bytes.
521        let (manifest, config_bytes, resolved_manifest_bytes) = self
522            .parse_and_resolve_manifest(&manifest_bytes, config_bytes, oci_ref)
523            .await?;
524
525        // Step 3: Parse config.
526        let (image_config, diff_ids) = ImageConfig::parse(&config_bytes)?;
527
528        // Step 4: Get layer descriptors.
529        let layer_descriptors = self.extract_layer_digests(&manifest)?;
530
531        // OCI spec requires diff_ids and layer descriptors to have the same count.
532        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        // Give the receiver a chance to render the resolved state before the
566        // layer tasks begin flooding download events.
567        tokio::task::yield_now().await;
568
569        // Warn about duplicate layer digests — they can cause contention.
570        {
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        // Materialize the shared EROFS layer stage and requested compositions.
584        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        // Clean up compressed tarballs after all layer tasks complete.
607        // Deferred from per-task cleanup to avoid races with duplicate layer digests.
608        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        // Persist cached image metadata.
614        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    /// Fetch manifest and config from the registry.
658    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    /// Parse manifest, resolving multi-platform index if needed.
674    ///
675    /// Returns the manifest and the correct config bytes. For single-platform
676    /// manifests, the config bytes are passed through unchanged. For multi-platform
677    /// indexes, the platform-specific config bytes are fetched and returned.
678    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        // Try to detect media type from the JSON.
685        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            // Resolve platform-specific manifest and fetch its config.
691            self.resolve_platform_manifest(manifest_bytes, reference)
692                .await
693        } else {
694            Ok((manifest, config_bytes, manifest_bytes.to_vec()))
695        }
696    }
697
698    /// Resolve a platform-specific manifest from an OCI index.
699    ///
700    /// Returns the resolved manifest and its platform-specific config bytes.
701    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        // Find matching platform.
712        let mut best_match: Option<&oci_spec::image::Descriptor> = None;
713        let mut exact_variant = false;
714
715        for entry in manifests {
716            // Skip attestation manifests.
717            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            // OS must match.
727            if *platform.os() != self.platform.os {
728                continue;
729            }
730
731            // Architecture must match.
732            if *platform.architecture() != self.platform.arch {
733                continue;
734            }
735
736            // Check variant.
737            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        // Fetch the platform-specific manifest and config.
762        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    /// Extract layer digests and sizes from a parsed manifest.
781    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    /// Materialize per-layer EROFS images, then generate fsmeta + VMDK.
815    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        // Validate all diff_ids parse as digests before spawning layer tasks.
831        // diff_ids come from the remote config blob (untrusted input).
832        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        // Phase-level idempotency is target-aware. Per-layer EROFS images are
843        // the common verified input for both output representations.
844        //
845        // - flat-only + all layers valid: no layer work; the caller builds ext4.
846        // - layered + fsmeta + VMDK valid: no-op.
847        // - layered + fsmeta valid + VMDK missing: re-stitch VMDK only.
848        // - layered + fsmeta missing: rebuild metadata from cached EROFS layers,
849        //   downloading and materializing only genuinely absent layers.
850        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        // `force` deliberately rebuilds every artifact. Otherwise a valid
876        // EROFS layer is reused even when a manifest-specific output is absent.
877        // Layer trees and block maps needed by fsmeta are reconstructed from
878        // EROFS below instead of forcing the source blob through the pipeline.
879        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                    // Acquire per-layer flock to coordinate with concurrent pulls.
947                    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                    // Re-check after lock — another process may have materialized it.
963                    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                    // Verify the uncompressed digest matches the config's diff_id.
1031                    // This is the OCI content trust check — the diff_id is signed
1032                    // as part of the image config, so a tampered layer would be caught.
1033                    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                    // Write to a temp file, then atomic rename to the final path.
1058                    // This prevents partial files from being visible to concurrent readers.
1059                    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                    // Strip file data from the retained tree to reduce memory.
1085                    // Only directory structure and metadata are needed for fsmeta merge.
1086                    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                    // Tarball cleanup is deferred — with duplicate layer digests,
1097                    // another task may still need the same blob. Tarballs are cleaned
1098                    // up after all layer tasks complete.
1099                    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        // Wait for all layer tasks to complete. Collect trees + data maps.
1118        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        // Generate fsmeta + VMDK if not already cached.
1126        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        // Acquire flock for fsmeta/VMDK generation.
1141        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        // Re-check after lock acquisition.
1154        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        // Newly written layers already carry their data-stripped tree and
1162        // block map. For cache hits, reconstruct the same inputs directly
1163        // from EROFS. This is what lets a new image reuse shared layers even
1164        // when it needs a new manifest-specific fsmeta artifact.
1165        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        // Merge layer trees with provenance tracking.
1180        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        // Generate fsmeta and VMDK.
1188        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        // Collect per-layer EROFS paths for the VMDK extents.
1194        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            // Write fsmeta.
1210            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            // Write VMDK descriptor.
1229            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    /// Re-stitch the VMDK descriptor from an existing fsmeta + layer EROFS files.
1262    ///
1263    /// Called when fsmeta and all layer EROFSes are present but only the VMDK
1264    /// descriptor is missing (e.g. the user deleted it manually, or a previous
1265    /// pull was interrupted between fsmeta rename and VMDK rename).
1266    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        // Re-check under lock: a concurrent pull may have regenerated VMDK,
1288        // or the fsmeta may have been evicted while we waited.
1289        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    // NOTE: materialize_flat_image was removed — replaced by fsmeta + VMDK generation
1350    // in materialize_layer_stage().
1351}
1352
1353//--------------------------------------------------------------------------------------------------
1354// Trait Implementations
1355//--------------------------------------------------------------------------------------------------
1356
1357impl<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
1394//--------------------------------------------------------------------------------------------------
1395// Functions: Helpers
1396//--------------------------------------------------------------------------------------------------
1397
1398/// Detect the media type of a manifest from its JSON content.
1399fn detect_manifest_media_type(bytes: &[u8]) -> String {
1400    // Try to parse the mediaType field from JSON.
1401    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        // Heuristic: if it has "manifests" array, it's an index.
1407        if v.get("manifests").is_some() {
1408            return "application/vnd.oci.image.index.v1+json".to_string();
1409        }
1410
1411        // If it has "layers" array, it's an image manifest.
1412        if v.get("layers").is_some() {
1413            return "application/vnd.oci.image.manifest.v1+json".to_string();
1414        }
1415    }
1416
1417    // Default to OCI image manifest.
1418    "application/vnd.oci.image.manifest.v1+json".to_string()
1419}
1420
1421/// Resolve the best matching platform-specific manifest digest.
1422pub(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
1459/// Build a pull result from cached image metadata.
1460fn 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
1476/// Assemble the metadata trees and block maps needed by fsmeta in layer order.
1477///
1478/// A freshly materialized layer supplies these values directly. Cache hits are
1479/// reconstructed from the verified EROFS artifact, so manifest-specific output
1480/// generation never requires the discarded registry tarball.
1481fn 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
1548/// Read only the metadata and regular-file block placement needed by fsmeta.
1549fn 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                    // Fsmeta takes the logical size and source blocks from
1575                    // ErofsDataMap, so retaining file contents is unnecessary.
1576                    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    // Check that all per-layer EROFS images exist.
1645    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//--------------------------------------------------------------------------------------------------
1872// Tests
1873//--------------------------------------------------------------------------------------------------
1874
1875#[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        // Create layers but no fsmeta/VMDK.
2234        let metadata = write_cached_image_fixture(&cache, &reference, &[false, false]);
2235        let manifest_digest = parse_digest(&metadata.manifest_digest);
2236        // Manually create layer files without fsmeta/VMDK.
2237        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        // Delete fsmeta/VMDK if they were created by the fixture.
2242        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        // Create EROFS files keyed by diff_id for cache hit detection.
2630        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        // Create fsmeta + VMDK when all layers are present (fsmerge pipeline).
2639        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        // Should build without error — we can't inspect ClientConfig directly,
2707        // but we verify it doesn't panic or fail.
2708        super::Registry::builder(Platform::default(), cache)
2709            .add_insecure_registries(vec!["localhost:5000".into()])
2710            .build()
2711            .unwrap();
2712    }
2713
2714    /// Generate a self-signed CA certificate and return PEM bytes.
2715    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    /// Helper to extract the error from a builder result.
2735    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        // Registry::new() should succeed just like builder().build()
2796        super::Registry::new(Platform::default(), cache).unwrap();
2797    }
2798}