Skip to main content

microsandbox_image/cache/
store.rs

1//! Global on-disk image and layer cache.
2
3use std::io::Read;
4use std::path::{Path, PathBuf};
5
6use oci_client::Reference;
7use serde::{Deserialize, Serialize};
8use sha2::{Digest as Sha2Digest, Sha256};
9
10use super::lock;
11use crate::{
12    config::ImageConfig,
13    digest::Digest,
14    erofs::ErofsReader,
15    error::{ImageError, ImageResult},
16    progress::{PullProgress, PullProgressSender},
17};
18
19//--------------------------------------------------------------------------------------------------
20// Constants
21//--------------------------------------------------------------------------------------------------
22
23/// Subdirectory for per-layer EROFS images (keyed by diff_id).
24const LAYERS_DIR: &str = "layers";
25
26/// Subdirectory for fsmeta EROFS images (keyed by manifest digest).
27const FSMETA_DIR: &str = "fsmeta";
28
29/// Subdirectory for VMDK descriptors (keyed by manifest digest).
30const VMDK_DIR: &str = "vmdk";
31
32/// Root directory for reusable flat ext4 artifacts.
33const FLAT_DIR: &str = "flat";
34const FLAT_REFS_DIR: &str = "refs";
35const FLAT_BLOBS_DIR: &str = "blobs";
36const FLAT_LOCKS_DIR: &str = "locks";
37
38/// Subdirectory for cached manifest + config metadata.
39const MANIFESTS_DIR: &str = "manifests";
40
41/// Subdirectory for transient staging (downloads, work dirs).
42const TMP_DIR: &str = "tmp";
43
44/// EROFS images are emitted in 4 KiB filesystem blocks.
45const EROFS_ALIGNMENT_BYTES: u64 = 4096;
46
47//--------------------------------------------------------------------------------------------------
48// Types
49//--------------------------------------------------------------------------------------------------
50
51/// Whether descriptor generation replaces an existing cached file.
52#[derive(Clone, Copy, PartialEq, Eq)]
53pub(crate) enum VmdkWriteMode {
54    /// Repair paths even when a descriptor already exists.
55    ReplaceExisting,
56    /// Recheck under the lock and preserve an existing descriptor.
57    KeepExisting,
58}
59
60/// On-disk global cache for OCI layers and EROFS images.
61///
62/// Layout:
63/// ```text
64/// ~/.microsandbox/cache/manifests/<sha256-of-ref>.json       # manifest + config metadata
65/// ~/.microsandbox/cache/tmp/<blob>.part                      # partial downloads
66/// ~/.microsandbox/cache/tmp/<blob>.download.lock             # download flock files
67/// ~/.microsandbox/cache/tmp/<blob>.work/                     # materialization work dirs
68/// ~/.microsandbox/cache/layers/<diff_id_safe>.erofs          # per-layer EROFS
69/// ~/.microsandbox/cache/layers/<diff_id_safe>.erofs.lock     # materialization flock
70/// ~/.microsandbox/cache/fsmeta/<manifest_safe>.erofs         # fsmeta EROFS (fsmerge metadata)
71/// ~/.microsandbox/cache/fsmeta/<manifest_safe>.erofs.lock    # materialization flock
72/// ~/.microsandbox/cache/vmdk/<manifest_safe>.vmdk            # VMDK descriptor
73/// ~/.microsandbox/cache/vmdk/<manifest_safe>.vmdk.lock       # materialization flock
74/// ```
75#[derive(Clone)]
76pub struct GlobalCache {
77    /// An explicitly scoped operation retains all admitted entries through catalog publication.
78    pub(super) operation: Option<
79        std::sync::Arc<
80            std::sync::Mutex<
81                std::collections::BTreeMap<PathBuf, crate::storage_lease::StorageLease>,
82            >,
83        >,
84    >,
85    /// Root of the layer EROFS cache (`~/.microsandbox/cache/layers/`).
86    layers_dir: PathBuf,
87
88    /// Root of the fsmeta EROFS cache (`~/.microsandbox/cache/fsmeta/`).
89    fsmeta_dir: PathBuf,
90
91    /// Root of the VMDK descriptor cache (`~/.microsandbox/cache/vmdk/`).
92    vmdk_dir: PathBuf,
93
94    /// Manifest-keyed references to immutable flat artifacts.
95    flat_refs_dir: PathBuf,
96
97    /// Content-addressed immutable raw ext4 artifacts.
98    flat_blobs_dir: PathBuf,
99
100    /// Per-derivation materialization locks.
101    flat_locks_dir: PathBuf,
102
103    /// Root of the manifest metadata cache (`~/.microsandbox/cache/manifests/`).
104    manifests_dir: PathBuf,
105
106    /// Root of the transient staging area (`~/.microsandbox/cache/tmp/`).
107    tmp_dir: PathBuf,
108}
109
110/// Cached metadata for a pulled image reference.
111#[derive(Debug, Clone, Serialize, Deserialize)]
112pub struct CachedImageMetadata {
113    /// Content-addressable digest of the resolved manifest.
114    pub manifest_digest: String,
115    /// Content-addressable digest of the config blob.
116    pub config_digest: String,
117    /// Raw resolved image manifest JSON.
118    pub raw_manifest_json: String,
119    /// Raw image config JSON.
120    pub raw_config_json: String,
121    /// Parsed OCI image configuration.
122    pub config: ImageConfig,
123    /// Layer metadata in bottom-to-top order.
124    pub layers: Vec<CachedLayerMetadata>,
125}
126
127/// Cached metadata for a single layer descriptor.
128#[derive(Debug, Clone, Serialize, Deserialize)]
129pub struct CachedLayerMetadata {
130    /// Compressed layer digest from the manifest (blob digest).
131    pub digest: String,
132    /// OCI media type of the layer blob.
133    pub media_type: Option<String>,
134    /// Compressed blob size in bytes.
135    pub size_bytes: Option<u64>,
136    /// Uncompressed diff ID from the image config.
137    pub diff_id: String,
138}
139
140/// Manifest-keyed reference to one validated immutable flat rootfs artifact.
141#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
142pub struct FlatRootfsRef {
143    /// Reference schema version.
144    pub schema: u32,
145    /// Resolved OCI manifest digest used as the requested input.
146    pub manifest_digest: String,
147    /// Complete derivation digest including platform and materializer profile.
148    pub derivation_digest: String,
149    /// SHA-256 digest of the validated raw ext4 bytes.
150    pub artifact_digest: String,
151    /// Pure-Rust materializer ABI.
152    pub materializer_abi: u32,
153    /// Deterministic ext4 UUID as lowercase hexadecimal.
154    pub uuid: String,
155    /// Logical sparse image size.
156    pub virtual_size_bytes: u64,
157    /// Unique inode count in the materialized rootfs.
158    pub inode_count: u64,
159    /// Unique regular-file content bytes.
160    pub content_bytes: u64,
161}
162
163//--------------------------------------------------------------------------------------------------
164// Methods
165//--------------------------------------------------------------------------------------------------
166
167impl GlobalCache {
168    /// Create a new GlobalCache using the provided cache directory.
169    ///
170    /// Creates all subdirectories if they don't exist.
171    pub fn new(cache_dir: &Path) -> ImageResult<Self> {
172        let layers_dir = cache_dir.join(LAYERS_DIR);
173        let fsmeta_dir = cache_dir.join(FSMETA_DIR);
174        let vmdk_dir = cache_dir.join(VMDK_DIR);
175        let flat_dir = cache_dir.join(FLAT_DIR);
176        let flat_refs_dir = flat_dir.join(FLAT_REFS_DIR);
177        let flat_blobs_dir = flat_dir.join(FLAT_BLOBS_DIR);
178        let flat_locks_dir = flat_dir.join(FLAT_LOCKS_DIR);
179        let manifests_dir = cache_dir.join(MANIFESTS_DIR);
180        let tmp_dir = cache_dir.join(TMP_DIR);
181
182        for dir in [
183            &layers_dir,
184            &fsmeta_dir,
185            &vmdk_dir,
186            &flat_refs_dir,
187            &flat_blobs_dir,
188            &flat_locks_dir,
189            &manifests_dir,
190            &tmp_dir,
191        ] {
192            std::fs::create_dir_all(dir).map_err(|e| ImageError::Cache {
193                path: dir.clone(),
194                source: e,
195            })?;
196        }
197
198        Ok(Self {
199            operation: None,
200            layers_dir,
201            fsmeta_dir,
202            vmdk_dir,
203            flat_refs_dir,
204            flat_blobs_dir,
205            flat_locks_dir,
206            manifests_dir,
207            tmp_dir,
208        })
209    }
210
211    /// Create a new GlobalCache using async filesystem operations.
212    pub async fn new_async(cache_dir: &Path) -> ImageResult<Self> {
213        let layers_dir = cache_dir.join(LAYERS_DIR);
214        let fsmeta_dir = cache_dir.join(FSMETA_DIR);
215        let vmdk_dir = cache_dir.join(VMDK_DIR);
216        let flat_dir = cache_dir.join(FLAT_DIR);
217        let flat_refs_dir = flat_dir.join(FLAT_REFS_DIR);
218        let flat_blobs_dir = flat_dir.join(FLAT_BLOBS_DIR);
219        let flat_locks_dir = flat_dir.join(FLAT_LOCKS_DIR);
220        let manifests_dir = cache_dir.join(MANIFESTS_DIR);
221        let tmp_dir = cache_dir.join(TMP_DIR);
222
223        for dir in [
224            &layers_dir,
225            &fsmeta_dir,
226            &vmdk_dir,
227            &flat_refs_dir,
228            &flat_blobs_dir,
229            &flat_locks_dir,
230            &manifests_dir,
231            &tmp_dir,
232        ] {
233            tokio::fs::create_dir_all(dir)
234                .await
235                .map_err(|e| ImageError::Cache {
236                    path: dir.clone(),
237                    source: e,
238                })?;
239        }
240
241        Ok(Self {
242            operation: None,
243            layers_dir,
244            fsmeta_dir,
245            vmdk_dir,
246            flat_refs_dir,
247            flat_blobs_dir,
248            flat_locks_dir,
249            manifests_dir,
250            tmp_dir,
251        })
252    }
253
254    // ── Layer EROFS paths (keyed by diff_id) ─────────────────────────
255
256    /// Root layer EROFS cache directory.
257    pub fn layers_dir(&self) -> &Path {
258        &self.layers_dir
259    }
260
261    /// Path to the per-layer EROFS image for a given diff_id.
262    pub fn layer_erofs_path(&self, diff_id: &Digest) -> PathBuf {
263        self.layers_dir
264            .join(format!("{}.erofs", diff_id.to_path_safe()))
265    }
266
267    /// Path to the materialization lock for a layer EROFS image.
268    pub fn layer_erofs_lock_path(&self, diff_id: &Digest) -> PathBuf {
269        self.layers_dir
270            .join(format!("{}.erofs.lock", diff_id.to_path_safe()))
271    }
272
273    /// Check if a layer EROFS image exists.
274    pub fn is_layer_materialized(&self, diff_id: &Digest) -> bool {
275        is_valid_erofs_artifact(&self.layer_erofs_path(diff_id))
276    }
277
278    /// Check if all given layer diff_ids have materialized EROFS images.
279    pub fn all_layers_materialized(&self, diff_ids: &[Digest]) -> bool {
280        diff_ids.iter().all(|d| self.is_layer_materialized(d))
281    }
282
283    // ── fsmeta EROFS paths (keyed by manifest digest) ─────────────────
284
285    /// Root fsmeta EROFS cache directory.
286    pub fn fsmeta_dir(&self) -> &Path {
287        &self.fsmeta_dir
288    }
289
290    /// Path to the fsmeta EROFS image for a given manifest digest.
291    pub fn fsmeta_erofs_path(&self, manifest_digest: &Digest) -> PathBuf {
292        self.fsmeta_dir
293            .join(format!("{}.erofs", manifest_digest.to_path_safe()))
294    }
295
296    /// Path to the materialization lock for a fsmeta EROFS image.
297    pub fn fsmeta_erofs_lock_path(&self, manifest_digest: &Digest) -> PathBuf {
298        self.fsmeta_dir
299            .join(format!("{}.erofs.lock", manifest_digest.to_path_safe()))
300    }
301
302    /// Check if a fsmeta EROFS image exists.
303    pub fn is_fsmeta_materialized(&self, manifest_digest: &Digest) -> bool {
304        is_valid_erofs_artifact(&self.fsmeta_erofs_path(manifest_digest))
305    }
306
307    // ── VMDK descriptor paths (keyed by manifest digest) ────────────
308
309    /// Root VMDK cache directory.
310    pub fn vmdk_dir(&self) -> &Path {
311        &self.vmdk_dir
312    }
313
314    /// Path to the VMDK descriptor for a given manifest digest.
315    pub fn vmdk_path(&self, manifest_digest: &Digest) -> PathBuf {
316        self.vmdk_dir
317            .join(format!("{}.vmdk", manifest_digest.to_path_safe()))
318    }
319
320    /// Path to the materialization lock for a VMDK descriptor.
321    pub fn vmdk_lock_path(&self, manifest_digest: &Digest) -> PathBuf {
322        self.vmdk_dir
323            .join(format!("{}.vmdk.lock", manifest_digest.to_path_safe()))
324    }
325
326    /// Check if a VMDK descriptor exists for a given manifest digest.
327    pub fn is_vmdk_materialized(&self, manifest_digest: &Digest) -> bool {
328        self.vmdk_path(manifest_digest).exists()
329    }
330
331    /// Rewrite the VMDK descriptor for a manifest digest so its extents are
332    /// this cache's fsmeta and layer EROFS files.
333    ///
334    /// Descriptors reference extents by absolute path, so one copied from
335    /// another cache must be regenerated rather than reused.
336    pub fn rewrite_vmdk(
337        &self,
338        manifest_digest: &Digest,
339        layer_diff_ids: &[Digest],
340    ) -> ImageResult<()> {
341        self.write_vmdk(
342            manifest_digest,
343            layer_diff_ids,
344            VmdkWriteMode::ReplaceExisting,
345            None,
346        )
347    }
348
349    /// Write a descriptor under the image materialization lock.
350    pub(crate) fn write_vmdk(
351        &self,
352        manifest_digest: &Digest,
353        layer_diff_ids: &[Digest],
354        mode: VmdkWriteMode,
355        progress: Option<&PullProgressSender>,
356    ) -> ImageResult<()> {
357        let fsmeta = self.fsmeta_erofs_path(manifest_digest);
358        let layers: Vec<PathBuf> = layer_diff_ids
359            .iter()
360            .map(|diff_id| self.layer_erofs_path(diff_id))
361            .collect();
362        let mut extents: Vec<&Path> = vec![&fsmeta];
363        extents.extend(layers.iter().map(PathBuf::as_path));
364
365        // Serialize with image materialization, which writes the VMDK under this lock.
366        let lock = lock::open_lock_file(&self.fsmeta_erofs_lock_path(manifest_digest))?;
367        lock::lock_exclusive(&lock)?;
368        let _unlock = scopeguard::guard(lock, |file| {
369            let _ = lock::flock_unlock(&file);
370        });
371
372        let vmdk = self.vmdk_path(manifest_digest);
373        if mode == VmdkWriteMode::KeepExisting {
374            // A concurrent pull may have regenerated the descriptor while we waited.
375            if vmdk.exists() {
376                return Ok(());
377            }
378            if !is_valid_erofs_artifact(&fsmeta) {
379                return Err(ImageError::Materialize {
380                    digest: manifest_digest.to_string(),
381                    message: "fsmeta vanished while waiting for VMDK regen lock".into(),
382                    source: None,
383                });
384            }
385        }
386
387        let work_dir = self.work_dir(manifest_digest);
388        std::fs::create_dir_all(&work_dir).map_err(|source| ImageError::Cache {
389            path: work_dir.clone(),
390            source,
391        })?;
392        let _work_guard = scopeguard::guard((), |_| {
393            let _ = std::fs::remove_dir_all(&work_dir);
394        });
395        if let Some(progress) = progress {
396            progress.send(PullProgress::StitchWritingVmdk);
397        }
398
399        let temp = work_dir.join("rootfs.vmdk");
400        crate::stitch::write_vmdk_descriptor(&temp, &extents).map_err(|source| match mode {
401            VmdkWriteMode::ReplaceExisting => ImageError::Cache {
402                path: vmdk.clone(),
403                source,
404            },
405            VmdkWriteMode::KeepExisting => ImageError::Materialize {
406                digest: manifest_digest.to_string(),
407                message: format!("VMDK write failed: {source}"),
408                source: None,
409            },
410        })?;
411        std::fs::rename(&temp, &vmdk).map_err(|source| ImageError::Cache { path: vmdk, source })?;
412
413        if let Some(progress) = progress {
414            progress.send(PullProgress::StitchComplete);
415        }
416
417        Ok(())
418    }
419
420    // ── Flat ext4 artifact paths (manifest ref → content blob) ───────
421
422    /// Path to the manifest-keyed flat rootfs reference.
423    pub fn flat_ref_path(&self, manifest_digest: &Digest) -> PathBuf {
424        self.flat_refs_dir
425            .join(format!("{}.json", manifest_digest.to_path_safe()))
426    }
427
428    /// Path to the immutable content-addressed raw ext4 artifact.
429    pub fn flat_blob_path(&self, artifact_digest: &Digest) -> PathBuf {
430        self.flat_blobs_dir
431            .join(format!("{}.raw", artifact_digest.to_path_safe()))
432    }
433
434    /// Path to the per-derivation flat materialization lock.
435    pub fn flat_lock_path(&self, derivation_digest: &Digest) -> PathBuf {
436        self.flat_locks_dir
437            .join(format!("{}.lock", derivation_digest.to_path_safe()))
438    }
439
440    /// Same-filesystem work directory for one flat-rootfs derivation.
441    pub fn flat_work_dir(&self, derivation_digest: &Digest) -> PathBuf {
442        self.tmp_dir
443            .join(format!("{}.flat.work", derivation_digest.to_path_safe()))
444    }
445
446    /// Publish a synchronized candidate as an immutable content-addressed blob.
447    pub fn publish_flat_blob(
448        &self,
449        candidate: &Path,
450        artifact_digest: &Digest,
451        expected_size: u64,
452    ) -> ImageResult<PathBuf> {
453        let destination = self.flat_blob_path(artifact_digest);
454        if let Ok(metadata) = std::fs::metadata(&destination) {
455            if metadata.len() != expected_size {
456                return Err(ImageError::Cache {
457                    path: destination,
458                    source: std::io::Error::new(
459                        std::io::ErrorKind::InvalidData,
460                        "content-addressed flat blob has an unexpected size",
461                    ),
462                });
463            }
464            if flat_blob_matches_digest(&destination, artifact_digest)? {
465                let _ = std::fs::remove_file(candidate);
466                return Ok(destination);
467            }
468
469            // A content-addressed name must never retain different bytes. The
470            // candidate has already been synchronized and validated, while a
471            // missing destination makes every existing ref safely miss after
472            // a crash between removal and rename.
473            std::fs::remove_file(&destination).map_err(|source| ImageError::Cache {
474                path: destination.clone(),
475                source,
476            })?;
477        }
478        std::fs::rename(candidate, &destination).map_err(|source| ImageError::Cache {
479            path: destination.clone(),
480            source,
481        })?;
482        sync_directory(&self.flat_blobs_dir)?;
483        Ok(destination)
484    }
485
486    /// Read and validate the manifest-keyed flat rootfs reference.
487    pub fn read_flat_ref(&self, manifest_digest: &Digest) -> ImageResult<Option<FlatRootfsRef>> {
488        let path = self.flat_ref_path(manifest_digest);
489        let data = match std::fs::read_to_string(&path) {
490            Ok(data) => data,
491            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
492            Err(source) => return Err(ImageError::Cache { path, source }),
493        };
494        let reference = match serde_json::from_str::<FlatRootfsRef>(&data) {
495            Ok(reference)
496                if reference.schema == 1
497                    && reference.manifest_digest == manifest_digest.to_string() =>
498            {
499                reference
500            }
501            Ok(_) => return Ok(None),
502            Err(error) => {
503                tracing::warn!(path = %path.display(), %error, "corrupt flat rootfs ref, ignoring");
504                return Ok(None);
505            }
506        };
507        let artifact_digest = match reference.artifact_digest.parse::<Digest>() {
508            Ok(digest) => digest,
509            Err(_) => return Ok(None),
510        };
511        let blob_path = self.flat_blob_path(&artifact_digest);
512        match std::fs::metadata(&blob_path) {
513            Ok(metadata) if metadata.len() == reference.virtual_size_bytes => Ok(Some(reference)),
514            Ok(_) | Err(_) => Ok(None),
515        }
516    }
517
518    /// Atomically replace a flat rootfs reference after its immutable blob is durable.
519    pub fn write_flat_ref(
520        &self,
521        manifest_digest: &Digest,
522        reference: &FlatRootfsRef,
523    ) -> ImageResult<()> {
524        let path = self.flat_ref_path(manifest_digest);
525        let temp_path = path.with_extension("json.part");
526        let payload = serde_json::to_vec_pretty(reference).map_err(|error| {
527            ImageError::ConfigParse(format!("failed to serialize flat rootfs ref: {error}"))
528        })?;
529        let mut temp = std::fs::File::create(&temp_path).map_err(|source| ImageError::Cache {
530            path: temp_path.clone(),
531            source,
532        })?;
533        use std::io::Write;
534        temp.write_all(&payload)
535            .map_err(|source| ImageError::Cache {
536                path: temp_path.clone(),
537                source,
538            })?;
539        temp.sync_all().map_err(|source| ImageError::Cache {
540            path: temp_path.clone(),
541            source,
542        })?;
543        std::fs::rename(&temp_path, &path).map_err(|source| ImageError::Cache {
544            path: path.clone(),
545            source,
546        })?;
547        sync_directory(&self.flat_refs_dir)?;
548        Ok(())
549    }
550
551    // ── Staging/tmp paths (downloads, work dirs) ─────────────────────
552
553    /// Root staging directory.
554    pub fn tmp_dir(&self) -> &Path {
555        &self.tmp_dir
556    }
557
558    /// Path to the partial download file for a blob.
559    pub fn part_path(&self, blob_digest: &Digest) -> PathBuf {
560        self.tmp_dir
561            .join(format!("{}.part", blob_digest.to_path_safe()))
562    }
563
564    /// Path to the download lock file for a blob.
565    pub fn download_lock_path(&self, blob_digest: &Digest) -> PathBuf {
566        self.tmp_dir
567            .join(format!("{}.download.lock", blob_digest.to_path_safe()))
568    }
569
570    /// Path to the materialization work directory for an EROFS build.
571    pub fn work_dir(&self, key: &Digest) -> PathBuf {
572        self.tmp_dir.join(format!("{}.work", key.to_path_safe()))
573    }
574
575    // ── Manifest metadata cache ──────────────────────────────────────
576
577    /// Root manifest metadata directory.
578    pub fn manifests_dir(&self) -> &Path {
579        &self.manifests_dir
580    }
581
582    /// Path to the pull lock file for an image reference.
583    pub fn image_lock_path(&self, reference: &Reference) -> PathBuf {
584        self.manifests_dir
585            .join(format!("{}.lock", image_cache_key(reference)))
586    }
587
588    /// Read cached metadata for an image reference.
589    pub fn read_image_metadata(
590        &self,
591        reference: &Reference,
592    ) -> ImageResult<Option<CachedImageMetadata>> {
593        self.read_image_metadata_path(&self.image_metadata_path(reference))
594    }
595
596    /// Resolve one mutable metadata entry and admit its exact dependencies under the
597    /// publication gate. Operation scopes retain those pins after this method returns.
598    pub(crate) fn read_image_metadata_path(
599        &self,
600        path: &Path,
601    ) -> ImageResult<Option<CachedImageMetadata>> {
602        self.read_image_metadata_matching(path, None)
603    }
604
605    fn read_image_metadata_matching(
606        &self,
607        path: &Path,
608        manifest_digest: Option<&str>,
609    ) -> ImageResult<Option<CachedImageMetadata>> {
610        // Scanning by digest must release each nonmatching entry immediately. Retaining
611        // every scanned image in the caller's operation can exhaust its descriptor limit.
612        let _entry = crate::storage_lease::StorageLease::shared(path)?;
613        let _gate =
614            crate::storage_lease::StorageLease::shared(&path.with_extension("publication"))?;
615        let data = match std::fs::read_to_string(path) {
616            Ok(data) => data,
617            Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
618            Err(source) => {
619                return Err(ImageError::Cache {
620                    path: path.to_path_buf(),
621                    source,
622                });
623            }
624        };
625        let metadata = parse_cached_image_metadata(path, &data)?;
626        if let Some(metadata) = &metadata {
627            if manifest_digest.is_some_and(|digest| digest != metadata.manifest_digest) {
628                return Ok(None);
629            }
630            let Ok(mut paths) = self.metadata_paths(metadata) else {
631                return Ok(None);
632            };
633            paths.push(path.to_path_buf());
634            let _dependencies = self.lease_paths(paths)?;
635        }
636        Ok(metadata)
637    }
638
639    /// Read and admit metadata without blocking the async executor on locks.
640    pub async fn read_image_metadata_async(
641        &self,
642        reference: &Reference,
643    ) -> ImageResult<Option<CachedImageMetadata>> {
644        self.read_image_metadata_path_async(self.image_metadata_path(reference))
645            .await
646    }
647
648    pub(crate) async fn read_image_metadata_path_async(
649        &self,
650        path: PathBuf,
651    ) -> ImageResult<Option<CachedImageMetadata>> {
652        let cache = self.clone();
653        tokio::task::spawn_blocking(move || cache.read_image_metadata_path(&path))
654            .await
655            .map_err(std::io::Error::other)?
656    }
657
658    pub(crate) async fn read_image_metadata_matching_async(
659        &self,
660        path: PathBuf,
661        manifest_digest: String,
662    ) -> ImageResult<Option<CachedImageMetadata>> {
663        let cache = self.clone();
664        tokio::task::spawn_blocking(move || {
665            cache.read_image_metadata_matching(&path, Some(&manifest_digest))
666        })
667        .await
668        .map_err(std::io::Error::other)?
669    }
670
671    /// Publish complete metadata under a short exclusive gate. Lifetime pins use a
672    /// different lock, allowing existing readers to finish using an older generation.
673    pub(crate) fn write_image_metadata(
674        &self,
675        reference: &Reference,
676        metadata: &CachedImageMetadata,
677    ) -> ImageResult<()> {
678        use std::io::Write;
679        let path = self.image_metadata_path(reference);
680        let mut paths = self.metadata_paths(metadata)?;
681        paths.push(path.clone());
682        let _leases = self.lease_paths(paths)?;
683        let _gate =
684            crate::storage_lease::StorageLease::exclusive(&path.with_extension("publication"))?;
685        let temporary_path = path.with_extension("json.part");
686        let mut temporary = std::fs::File::create(&temporary_path)?;
687        serde_json::to_writer(&mut temporary, metadata)
688            .map_err(|error| ImageError::ConfigParse(error.to_string()))?;
689        temporary.flush()?;
690        temporary.sync_all()?;
691        std::fs::rename(&temporary_path, &path)?;
692        sync_directory(&self.manifests_dir)?;
693        Ok(())
694    }
695
696    /// Publish metadata from an owned worker so cancellation cannot release its gate
697    /// while the filesystem replacement is still in flight.
698    pub async fn write_image_metadata_async(
699        &self,
700        reference: &Reference,
701        metadata: &CachedImageMetadata,
702    ) -> ImageResult<()> {
703        let cache = self.clone();
704        let reference = reference.clone();
705        let metadata = metadata.clone();
706        tokio::task::spawn_blocking(move || cache.write_image_metadata(&reference, &metadata))
707            .await
708            .map_err(std::io::Error::other)?
709    }
710
711    /// Install archived metadata without replacing an equivalent local entry. Recheck under
712    /// the same publication gate used by image pulls, so a concurrent retag cannot be lost.
713    pub async fn install_image_metadata_if_compatible_async(
714        &self,
715        reference: &Reference,
716        archived_bytes: Vec<u8>,
717    ) -> ImageResult<()> {
718        let cache = self.clone();
719        let reference = reference.clone();
720        tokio::task::spawn_blocking(move || {
721            cache.install_image_metadata_if_compatible(&reference, &archived_bytes)
722        })
723        .await
724        .map_err(std::io::Error::other)?
725    }
726
727    fn install_image_metadata_if_compatible(
728        &self,
729        reference: &Reference,
730        archived_bytes: &[u8],
731    ) -> ImageResult<()> {
732        use std::io::Write;
733
734        crate::snapshot::manifest::reject_duplicate_json_keys(archived_bytes)
735            .map_err(|error| ImageError::ConfigParse(error.to_string()))?;
736        let metadata: CachedImageMetadata = serde_json::from_slice(archived_bytes)
737            .map_err(|error| ImageError::ConfigParse(error.to_string()))?;
738        let path = self.image_metadata_path(reference);
739        let mut paths = self.metadata_paths(&metadata)?;
740        paths.push(path.clone());
741        let _leases = self.lease_paths(paths)?;
742        let _gate =
743            crate::storage_lease::StorageLease::exclusive(&path.with_extension("publication"))?;
744
745        match std::fs::symlink_metadata(&path) {
746            Ok(existing) if existing.file_type().is_file() => {
747                let existing_bytes = std::fs::read(&path)?;
748                if image_metadata_json_equivalent(archived_bytes, &existing_bytes) {
749                    return Ok(());
750                }
751                return Err(ImageError::Cache {
752                    path,
753                    source: std::io::Error::new(
754                        std::io::ErrorKind::AlreadyExists,
755                        "cache target already exists with different content",
756                    ),
757                });
758            }
759            Ok(_) => {
760                return Err(ImageError::Cache {
761                    path,
762                    source: std::io::Error::new(
763                        std::io::ErrorKind::InvalidData,
764                        "cache target is not a regular file",
765                    ),
766                });
767            }
768            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
769            Err(source) => return Err(ImageError::Cache { path, source }),
770        }
771
772        // The gate also covers the existing writer's .json.part path. Preserve the archive's
773        // bytes on first install, then publish only after the complete file is durable.
774        let temporary_path = path.with_extension("json.part");
775        let mut temporary = std::fs::File::create(&temporary_path)?;
776        temporary.write_all(archived_bytes)?;
777        temporary.sync_all()?;
778        std::fs::rename(temporary_path, path)?;
779        sync_directory(&self.manifests_dir)?;
780        Ok(())
781    }
782
783    /// Delete cached metadata for an image reference.
784    pub fn delete_image_metadata(&self, reference: &Reference) -> ImageResult<()> {
785        let path = self.image_metadata_path(reference);
786        match std::fs::remove_file(&path) {
787            Ok(()) => Ok(()),
788            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
789            Err(e) => Err(ImageError::Cache { path, source: e }),
790        }
791    }
792
793    /// Delete cached metadata for an image reference using async filesystem I/O.
794    pub async fn delete_image_metadata_async(&self, reference: &Reference) -> ImageResult<()> {
795        let path = self.image_metadata_path(reference);
796        match tokio::fs::remove_file(&path).await {
797            Ok(()) => Ok(()),
798            Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
799            Err(e) => Err(ImageError::Cache { path, source: e }),
800        }
801    }
802
803    /// Path to the cached metadata file for an image reference.
804    pub fn image_metadata_path(&self, reference: &Reference) -> PathBuf {
805        self.manifests_dir
806            .join(format!("{}.json", image_cache_key(reference)))
807    }
808
809    // ── Blob cache paths ──────────────────────────────────────────────
810
811    /// Path to the cached compressed tarball for a layer blob.
812    pub fn tar_path(&self, digest: &Digest) -> PathBuf {
813        self.layers_dir
814            .join(format!("{}.tar.gz", digest.to_path_safe()))
815    }
816}
817
818//--------------------------------------------------------------------------------------------------
819// Functions
820//--------------------------------------------------------------------------------------------------
821
822fn image_cache_key(reference: &Reference) -> String {
823    let mut hasher = Sha256::new();
824    hasher.update(reference.to_string().as_bytes());
825    hex::encode(hasher.finalize())
826}
827
828fn image_metadata_json_equivalent(left: &[u8], right: &[u8]) -> bool {
829    if crate::snapshot::manifest::reject_duplicate_json_keys(left).is_err()
830        || crate::snapshot::manifest::reject_duplicate_json_keys(right).is_err()
831    {
832        return false;
833    }
834    matches!(
835        (
836            serde_json::from_slice::<serde_json::Value>(left),
837            serde_json::from_slice::<serde_json::Value>(right),
838        ),
839        (Ok(left), Ok(right)) if left == right
840    )
841}
842
843#[cfg(unix)]
844fn sync_directory(path: &Path) -> ImageResult<()> {
845    let directory = std::fs::File::open(path).map_err(|source| ImageError::Cache {
846        path: path.to_path_buf(),
847        source,
848    })?;
849    directory.sync_all().map_err(|source| ImageError::Cache {
850        path: path.to_path_buf(),
851        source,
852    })
853}
854
855#[cfg(not(unix))]
856fn sync_directory(_path: &Path) -> ImageResult<()> {
857    Ok(())
858}
859
860pub(crate) fn parse_cached_image_metadata(
861    path: &Path,
862    data: &str,
863) -> ImageResult<Option<CachedImageMetadata>> {
864    match serde_json::from_str::<CachedImageMetadata>(data) {
865        Ok(metadata) => Ok(Some(metadata)),
866        Err(e) => {
867            tracing::warn!(
868                path = %path.display(),
869                error = %e,
870                "corrupt image metadata cache, ignoring"
871            );
872            Ok(None)
873        }
874    }
875}
876
877pub(crate) fn is_valid_erofs_artifact(path: &Path) -> bool {
878    let Ok(meta) = std::fs::metadata(path) else {
879        return false;
880    };
881    let len = meta.len();
882    if !meta.is_file() || len == 0 || len % EROFS_ALIGNMENT_BYTES != 0 {
883        return false;
884    }
885
886    // Length alone accepts any aligned garbage as a cache hit. Parse the
887    // superblock and root inode so corrupt cached layers are re-materialized
888    // instead of failing later while composing a flat rootfs or booting a VM.
889    let Ok(file) = std::fs::File::open(path) else {
890        return false;
891    };
892    let Ok(mut reader) = ErofsReader::new(file) else {
893        return false;
894    };
895    reader.root_directory_metadata().is_ok()
896}
897
898pub(crate) async fn is_valid_erofs_artifact_async(path: &Path) -> bool {
899    let path = path.to_path_buf();
900    tokio::task::spawn_blocking(move || is_valid_erofs_artifact(&path))
901        .await
902        .unwrap_or(false)
903}
904
905fn flat_blob_matches_digest(path: &Path, expected: &Digest) -> ImageResult<bool> {
906    let mut file = std::fs::File::open(path).map_err(|source| ImageError::Cache {
907        path: path.to_path_buf(),
908        source,
909    })?;
910    let mut hasher = Sha256::new();
911    let mut buffer = [0u8; 1024 * 1024];
912    loop {
913        let read = file.read(&mut buffer).map_err(|source| ImageError::Cache {
914            path: path.to_path_buf(),
915            source,
916        })?;
917        if read == 0 {
918            break;
919        }
920        hasher.update(&buffer[..read]);
921    }
922    Ok(format!("sha256:{}", hex::encode(hasher.finalize())) == expected.to_string())
923}
924
925//--------------------------------------------------------------------------------------------------
926// Tests
927//--------------------------------------------------------------------------------------------------
928
929#[cfg(test)]
930mod tests {
931    use super::*;
932
933    fn digest(byte: char) -> Digest {
934        format!("sha256:{}", byte.to_string().repeat(64))
935            .parse()
936            .unwrap()
937    }
938
939    #[tokio::test]
940    async fn archived_metadata_rechecks_after_a_concurrent_retag() {
941        let directory = tempfile::tempdir().unwrap();
942        let cache = GlobalCache::new(directory.path()).unwrap().operation();
943        let reference: Reference = "example.com/test:latest".parse().unwrap();
944        let archived = CachedImageMetadata {
945            manifest_digest: digest('a').to_string(),
946            config_digest: digest('a').to_string(),
947            raw_manifest_json: "{}".into(),
948            raw_config_json: "{}".into(),
949            config: ImageConfig::default(),
950            layers: Vec::new(),
951        };
952        let archived_bytes = serde_json::to_vec(&archived).unwrap();
953        cache.write_image_metadata(&reference, &archived).unwrap();
954        let path = cache.image_metadata_path(&reference);
955        assert!(image_metadata_json_equivalent(
956            &archived_bytes,
957            &std::fs::read(&path).unwrap()
958        ));
959
960        // Model another image pull changing the mutable tag after archive preflight.
961        let mut retagged = archived.clone();
962        retagged.manifest_digest = digest('b').to_string();
963        cache.write_image_metadata(&reference, &retagged).unwrap();
964        let retagged_bytes = std::fs::read(&path).unwrap();
965
966        let error = cache
967            .install_image_metadata_if_compatible_async(&reference, archived_bytes)
968            .await
969            .unwrap_err();
970        assert!(
971            error
972                .to_string()
973                .contains("cache target already exists with different content")
974        );
975        assert_eq!(std::fs::read(path).unwrap(), retagged_bytes);
976    }
977
978    #[test]
979    fn flat_cache_separates_manifest_refs_from_content_blobs() {
980        let directory = tempfile::tempdir().unwrap();
981        let cache = GlobalCache::new(directory.path()).unwrap();
982        let manifest = digest('a');
983        let derivation = digest('b');
984        let artifact = digest('c');
985
986        assert!(cache.flat_ref_path(&manifest).ends_with(
987            "flat/refs/sha256_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.json"
988        ));
989        assert!(cache.flat_blob_path(&artifact).ends_with(
990            "flat/blobs/sha256_cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc.raw"
991        ));
992        assert!(cache.flat_lock_path(&derivation).ends_with(
993            "flat/locks/sha256_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb.lock"
994        ));
995
996        let blob_path = cache.flat_blob_path(&artifact);
997        let blob = std::fs::File::create(&blob_path).unwrap();
998        blob.set_len(4096).unwrap();
999        let reference = FlatRootfsRef {
1000            schema: 1,
1001            manifest_digest: manifest.to_string(),
1002            derivation_digest: derivation.to_string(),
1003            artifact_digest: artifact.to_string(),
1004            materializer_abi: 1,
1005            uuid: "00".repeat(16),
1006            virtual_size_bytes: 4096,
1007            inode_count: 2,
1008            content_bytes: 7,
1009        };
1010        cache.write_flat_ref(&manifest, &reference).unwrap();
1011
1012        assert_eq!(cache.read_flat_ref(&manifest).unwrap(), Some(reference));
1013    }
1014
1015    #[test]
1016    fn aligned_garbage_is_not_a_valid_erofs_cache_artifact() {
1017        let directory = tempfile::tempdir().unwrap();
1018        let path = directory.path().join("corrupt.erofs");
1019        let file = std::fs::File::create(&path).unwrap();
1020        file.set_len(EROFS_ALIGNMENT_BYTES).unwrap();
1021
1022        assert!(!is_valid_erofs_artifact(&path));
1023    }
1024
1025    #[test]
1026    fn flat_ref_must_name_the_manifest_that_indexes_it() {
1027        let directory = tempfile::tempdir().unwrap();
1028        let cache = GlobalCache::new(directory.path()).unwrap();
1029        let indexed_manifest = digest('a');
1030        let wrong_manifest = digest('b');
1031        let artifact = digest('c');
1032        std::fs::write(cache.flat_blob_path(&artifact), [0u8; 8]).unwrap();
1033        let reference = FlatRootfsRef {
1034            schema: 1,
1035            manifest_digest: wrong_manifest.to_string(),
1036            derivation_digest: digest('d').to_string(),
1037            artifact_digest: artifact.to_string(),
1038            materializer_abi: 1,
1039            uuid: "00".repeat(16),
1040            virtual_size_bytes: 8,
1041            inode_count: 2,
1042            content_bytes: 0,
1043        };
1044        cache.write_flat_ref(&indexed_manifest, &reference).unwrap();
1045
1046        assert_eq!(cache.read_flat_ref(&indexed_manifest).unwrap(), None);
1047    }
1048
1049    #[test]
1050    fn publishing_replaces_same_size_blob_with_wrong_content() {
1051        let directory = tempfile::tempdir().unwrap();
1052        let cache = GlobalCache::new(directory.path()).unwrap();
1053        let candidate = directory.path().join("candidate.raw");
1054        let expected_bytes = b"good";
1055        std::fs::write(&candidate, expected_bytes).unwrap();
1056        let expected: Digest = format!("sha256:{}", hex::encode(Sha256::digest(expected_bytes)))
1057            .parse()
1058            .unwrap();
1059        let destination = cache.flat_blob_path(&expected);
1060        std::fs::write(&destination, b"evil").unwrap();
1061
1062        cache
1063            .publish_flat_blob(&candidate, &expected, expected_bytes.len() as u64)
1064            .unwrap();
1065
1066        assert_eq!(std::fs::read(destination).unwrap(), expected_bytes);
1067    }
1068}