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