use std::io::Read;
use std::path::{Path, PathBuf};
use oci_client::Reference;
use serde::{Deserialize, Serialize};
use sha2::{Digest as Sha2Digest, Sha256};
use crate::{
config::ImageConfig,
digest::Digest,
erofs::ErofsReader,
error::{ImageError, ImageResult},
};
const LAYERS_DIR: &str = "layers";
const FSMETA_DIR: &str = "fsmeta";
const VMDK_DIR: &str = "vmdk";
const FLAT_DIR: &str = "flat";
const FLAT_REFS_DIR: &str = "refs";
const FLAT_BLOBS_DIR: &str = "blobs";
const FLAT_LOCKS_DIR: &str = "locks";
const MANIFESTS_DIR: &str = "manifests";
const TMP_DIR: &str = "tmp";
const EROFS_ALIGNMENT_BYTES: u64 = 4096;
#[derive(Clone)]
pub struct GlobalCache {
layers_dir: PathBuf,
fsmeta_dir: PathBuf,
vmdk_dir: PathBuf,
flat_refs_dir: PathBuf,
flat_blobs_dir: PathBuf,
flat_locks_dir: PathBuf,
manifests_dir: PathBuf,
tmp_dir: PathBuf,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CachedImageMetadata {
pub manifest_digest: String,
pub config_digest: String,
pub raw_manifest_json: String,
pub raw_config_json: String,
pub config: ImageConfig,
pub layers: Vec<CachedLayerMetadata>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CachedLayerMetadata {
pub digest: String,
pub media_type: Option<String>,
pub size_bytes: Option<u64>,
pub diff_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlatRootfsRef {
pub schema: u32,
pub manifest_digest: String,
pub derivation_digest: String,
pub artifact_digest: String,
pub materializer_abi: u32,
pub uuid: String,
pub virtual_size_bytes: u64,
pub inode_count: u64,
pub content_bytes: u64,
}
impl GlobalCache {
pub fn new(cache_dir: &Path) -> ImageResult<Self> {
let layers_dir = cache_dir.join(LAYERS_DIR);
let fsmeta_dir = cache_dir.join(FSMETA_DIR);
let vmdk_dir = cache_dir.join(VMDK_DIR);
let flat_dir = cache_dir.join(FLAT_DIR);
let flat_refs_dir = flat_dir.join(FLAT_REFS_DIR);
let flat_blobs_dir = flat_dir.join(FLAT_BLOBS_DIR);
let flat_locks_dir = flat_dir.join(FLAT_LOCKS_DIR);
let manifests_dir = cache_dir.join(MANIFESTS_DIR);
let tmp_dir = cache_dir.join(TMP_DIR);
for dir in [
&layers_dir,
&fsmeta_dir,
&vmdk_dir,
&flat_refs_dir,
&flat_blobs_dir,
&flat_locks_dir,
&manifests_dir,
&tmp_dir,
] {
std::fs::create_dir_all(dir).map_err(|e| ImageError::Cache {
path: dir.clone(),
source: e,
})?;
}
Ok(Self {
layers_dir,
fsmeta_dir,
vmdk_dir,
flat_refs_dir,
flat_blobs_dir,
flat_locks_dir,
manifests_dir,
tmp_dir,
})
}
pub async fn new_async(cache_dir: &Path) -> ImageResult<Self> {
let layers_dir = cache_dir.join(LAYERS_DIR);
let fsmeta_dir = cache_dir.join(FSMETA_DIR);
let vmdk_dir = cache_dir.join(VMDK_DIR);
let flat_dir = cache_dir.join(FLAT_DIR);
let flat_refs_dir = flat_dir.join(FLAT_REFS_DIR);
let flat_blobs_dir = flat_dir.join(FLAT_BLOBS_DIR);
let flat_locks_dir = flat_dir.join(FLAT_LOCKS_DIR);
let manifests_dir = cache_dir.join(MANIFESTS_DIR);
let tmp_dir = cache_dir.join(TMP_DIR);
for dir in [
&layers_dir,
&fsmeta_dir,
&vmdk_dir,
&flat_refs_dir,
&flat_blobs_dir,
&flat_locks_dir,
&manifests_dir,
&tmp_dir,
] {
tokio::fs::create_dir_all(dir)
.await
.map_err(|e| ImageError::Cache {
path: dir.clone(),
source: e,
})?;
}
Ok(Self {
layers_dir,
fsmeta_dir,
vmdk_dir,
flat_refs_dir,
flat_blobs_dir,
flat_locks_dir,
manifests_dir,
tmp_dir,
})
}
pub fn layers_dir(&self) -> &Path {
&self.layers_dir
}
pub fn layer_erofs_path(&self, diff_id: &Digest) -> PathBuf {
self.layers_dir
.join(format!("{}.erofs", diff_id.to_path_safe()))
}
pub fn layer_erofs_lock_path(&self, diff_id: &Digest) -> PathBuf {
self.layers_dir
.join(format!("{}.erofs.lock", diff_id.to_path_safe()))
}
pub fn is_layer_materialized(&self, diff_id: &Digest) -> bool {
is_valid_erofs_artifact(&self.layer_erofs_path(diff_id))
}
pub fn all_layers_materialized(&self, diff_ids: &[Digest]) -> bool {
diff_ids.iter().all(|d| self.is_layer_materialized(d))
}
pub fn fsmeta_dir(&self) -> &Path {
&self.fsmeta_dir
}
pub fn fsmeta_erofs_path(&self, manifest_digest: &Digest) -> PathBuf {
self.fsmeta_dir
.join(format!("{}.erofs", manifest_digest.to_path_safe()))
}
pub fn fsmeta_erofs_lock_path(&self, manifest_digest: &Digest) -> PathBuf {
self.fsmeta_dir
.join(format!("{}.erofs.lock", manifest_digest.to_path_safe()))
}
pub fn is_fsmeta_materialized(&self, manifest_digest: &Digest) -> bool {
is_valid_erofs_artifact(&self.fsmeta_erofs_path(manifest_digest))
}
pub fn vmdk_dir(&self) -> &Path {
&self.vmdk_dir
}
pub fn vmdk_path(&self, manifest_digest: &Digest) -> PathBuf {
self.vmdk_dir
.join(format!("{}.vmdk", manifest_digest.to_path_safe()))
}
pub fn vmdk_lock_path(&self, manifest_digest: &Digest) -> PathBuf {
self.vmdk_dir
.join(format!("{}.vmdk.lock", manifest_digest.to_path_safe()))
}
pub fn is_vmdk_materialized(&self, manifest_digest: &Digest) -> bool {
self.vmdk_path(manifest_digest).exists()
}
pub fn flat_ref_path(&self, manifest_digest: &Digest) -> PathBuf {
self.flat_refs_dir
.join(format!("{}.json", manifest_digest.to_path_safe()))
}
pub fn flat_blob_path(&self, artifact_digest: &Digest) -> PathBuf {
self.flat_blobs_dir
.join(format!("{}.raw", artifact_digest.to_path_safe()))
}
pub fn flat_lock_path(&self, derivation_digest: &Digest) -> PathBuf {
self.flat_locks_dir
.join(format!("{}.lock", derivation_digest.to_path_safe()))
}
pub fn flat_work_dir(&self, derivation_digest: &Digest) -> PathBuf {
self.tmp_dir
.join(format!("{}.flat.work", derivation_digest.to_path_safe()))
}
pub fn publish_flat_blob(
&self,
candidate: &Path,
artifact_digest: &Digest,
expected_size: u64,
) -> ImageResult<PathBuf> {
let destination = self.flat_blob_path(artifact_digest);
if let Ok(metadata) = std::fs::metadata(&destination) {
if metadata.len() != expected_size {
return Err(ImageError::Cache {
path: destination,
source: std::io::Error::new(
std::io::ErrorKind::InvalidData,
"content-addressed flat blob has an unexpected size",
),
});
}
if flat_blob_matches_digest(&destination, artifact_digest)? {
let _ = std::fs::remove_file(candidate);
return Ok(destination);
}
std::fs::remove_file(&destination).map_err(|source| ImageError::Cache {
path: destination.clone(),
source,
})?;
}
std::fs::rename(candidate, &destination).map_err(|source| ImageError::Cache {
path: destination.clone(),
source,
})?;
sync_directory(&self.flat_blobs_dir)?;
Ok(destination)
}
pub fn read_flat_ref(&self, manifest_digest: &Digest) -> ImageResult<Option<FlatRootfsRef>> {
let path = self.flat_ref_path(manifest_digest);
let data = match std::fs::read_to_string(&path) {
Ok(data) => data,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(source) => return Err(ImageError::Cache { path, source }),
};
let reference = match serde_json::from_str::<FlatRootfsRef>(&data) {
Ok(reference)
if reference.schema == 1
&& reference.manifest_digest == manifest_digest.to_string() =>
{
reference
}
Ok(_) => return Ok(None),
Err(error) => {
tracing::warn!(path = %path.display(), %error, "corrupt flat rootfs ref, ignoring");
return Ok(None);
}
};
let artifact_digest = match reference.artifact_digest.parse::<Digest>() {
Ok(digest) => digest,
Err(_) => return Ok(None),
};
let blob_path = self.flat_blob_path(&artifact_digest);
match std::fs::metadata(&blob_path) {
Ok(metadata) if metadata.len() == reference.virtual_size_bytes => Ok(Some(reference)),
Ok(_) | Err(_) => Ok(None),
}
}
pub fn write_flat_ref(
&self,
manifest_digest: &Digest,
reference: &FlatRootfsRef,
) -> ImageResult<()> {
let path = self.flat_ref_path(manifest_digest);
let temp_path = path.with_extension("json.part");
let payload = serde_json::to_vec_pretty(reference).map_err(|error| {
ImageError::ConfigParse(format!("failed to serialize flat rootfs ref: {error}"))
})?;
let mut temp = std::fs::File::create(&temp_path).map_err(|source| ImageError::Cache {
path: temp_path.clone(),
source,
})?;
use std::io::Write;
temp.write_all(&payload)
.map_err(|source| ImageError::Cache {
path: temp_path.clone(),
source,
})?;
temp.sync_all().map_err(|source| ImageError::Cache {
path: temp_path.clone(),
source,
})?;
std::fs::rename(&temp_path, &path).map_err(|source| ImageError::Cache {
path: path.clone(),
source,
})?;
sync_directory(&self.flat_refs_dir)?;
Ok(())
}
pub fn tmp_dir(&self) -> &Path {
&self.tmp_dir
}
pub fn part_path(&self, blob_digest: &Digest) -> PathBuf {
self.tmp_dir
.join(format!("{}.part", blob_digest.to_path_safe()))
}
pub fn download_lock_path(&self, blob_digest: &Digest) -> PathBuf {
self.tmp_dir
.join(format!("{}.download.lock", blob_digest.to_path_safe()))
}
pub fn work_dir(&self, key: &Digest) -> PathBuf {
self.tmp_dir.join(format!("{}.work", key.to_path_safe()))
}
pub fn manifests_dir(&self) -> &Path {
&self.manifests_dir
}
pub fn image_lock_path(&self, reference: &Reference) -> PathBuf {
self.manifests_dir
.join(format!("{}.lock", image_cache_key(reference)))
}
pub fn read_image_metadata(
&self,
reference: &Reference,
) -> ImageResult<Option<CachedImageMetadata>> {
let path = self.image_metadata_path(reference);
let data = match std::fs::read_to_string(&path) {
Ok(data) => data,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(ImageError::Cache { path, source: e }),
};
parse_cached_image_metadata(&path, &data)
}
pub async fn read_image_metadata_async(
&self,
reference: &Reference,
) -> ImageResult<Option<CachedImageMetadata>> {
let path = self.image_metadata_path(reference);
let data = match tokio::fs::read_to_string(&path).await {
Ok(data) => data,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(e) => return Err(ImageError::Cache { path, source: e }),
};
parse_cached_image_metadata(&path, &data)
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn write_image_metadata(
&self,
reference: &Reference,
metadata: &CachedImageMetadata,
) -> ImageResult<()> {
let path = self.image_metadata_path(reference);
let temp_path = path.with_extension("json.part");
let payload = serde_json::to_vec(metadata).map_err(|e| {
ImageError::ConfigParse(format!("failed to serialize cached image metadata: {e}"))
})?;
std::fs::write(&temp_path, payload).map_err(|e| ImageError::Cache {
path: temp_path.clone(),
source: e,
})?;
std::fs::rename(&temp_path, &path).map_err(|e| ImageError::Cache { path, source: e })?;
Ok(())
}
pub async fn write_image_metadata_async(
&self,
reference: &Reference,
metadata: &CachedImageMetadata,
) -> ImageResult<()> {
let path = self.image_metadata_path(reference);
let temp_path = path.with_extension("json.part");
let payload = serde_json::to_vec(metadata).map_err(|e| {
ImageError::ConfigParse(format!("failed to serialize cached image metadata: {e}"))
})?;
tokio::fs::write(&temp_path, payload)
.await
.map_err(|e| ImageError::Cache {
path: temp_path.clone(),
source: e,
})?;
tokio::fs::rename(&temp_path, &path)
.await
.map_err(|e| ImageError::Cache { path, source: e })?;
Ok(())
}
pub fn delete_image_metadata(&self, reference: &Reference) -> ImageResult<()> {
let path = self.image_metadata_path(reference);
match std::fs::remove_file(&path) {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(ImageError::Cache { path, source: e }),
}
}
pub async fn delete_image_metadata_async(&self, reference: &Reference) -> ImageResult<()> {
let path = self.image_metadata_path(reference);
match tokio::fs::remove_file(&path).await {
Ok(()) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()),
Err(e) => Err(ImageError::Cache { path, source: e }),
}
}
pub fn image_metadata_path(&self, reference: &Reference) -> PathBuf {
self.manifests_dir
.join(format!("{}.json", image_cache_key(reference)))
}
pub fn tar_path(&self, digest: &Digest) -> PathBuf {
self.layers_dir
.join(format!("{}.tar.gz", digest.to_path_safe()))
}
}
fn image_cache_key(reference: &Reference) -> String {
let mut hasher = Sha256::new();
hasher.update(reference.to_string().as_bytes());
hex::encode(hasher.finalize())
}
#[cfg(unix)]
fn sync_directory(path: &Path) -> ImageResult<()> {
let directory = std::fs::File::open(path).map_err(|source| ImageError::Cache {
path: path.to_path_buf(),
source,
})?;
directory.sync_all().map_err(|source| ImageError::Cache {
path: path.to_path_buf(),
source,
})
}
#[cfg(not(unix))]
fn sync_directory(_path: &Path) -> ImageResult<()> {
Ok(())
}
pub(crate) fn parse_cached_image_metadata(
path: &Path,
data: &str,
) -> ImageResult<Option<CachedImageMetadata>> {
match serde_json::from_str::<CachedImageMetadata>(data) {
Ok(metadata) => Ok(Some(metadata)),
Err(e) => {
tracing::warn!(
path = %path.display(),
error = %e,
"corrupt image metadata cache, ignoring"
);
Ok(None)
}
}
}
pub(crate) fn is_valid_erofs_artifact(path: &Path) -> bool {
let Ok(meta) = std::fs::metadata(path) else {
return false;
};
let len = meta.len();
if !meta.is_file() || len == 0 || len % EROFS_ALIGNMENT_BYTES != 0 {
return false;
}
let Ok(file) = std::fs::File::open(path) else {
return false;
};
let Ok(mut reader) = ErofsReader::new(file) else {
return false;
};
reader.root_directory_metadata().is_ok()
}
pub(crate) async fn is_valid_erofs_artifact_async(path: &Path) -> bool {
let path = path.to_path_buf();
tokio::task::spawn_blocking(move || is_valid_erofs_artifact(&path))
.await
.unwrap_or(false)
}
fn flat_blob_matches_digest(path: &Path, expected: &Digest) -> ImageResult<bool> {
let mut file = std::fs::File::open(path).map_err(|source| ImageError::Cache {
path: path.to_path_buf(),
source,
})?;
let mut hasher = Sha256::new();
let mut buffer = [0u8; 1024 * 1024];
loop {
let read = file.read(&mut buffer).map_err(|source| ImageError::Cache {
path: path.to_path_buf(),
source,
})?;
if read == 0 {
break;
}
hasher.update(&buffer[..read]);
}
Ok(format!("sha256:{}", hex::encode(hasher.finalize())) == expected.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
fn digest(byte: char) -> Digest {
format!("sha256:{}", byte.to_string().repeat(64))
.parse()
.unwrap()
}
#[test]
fn flat_cache_separates_manifest_refs_from_content_blobs() {
let directory = tempfile::tempdir().unwrap();
let cache = GlobalCache::new(directory.path()).unwrap();
let manifest = digest('a');
let derivation = digest('b');
let artifact = digest('c');
assert!(cache.flat_ref_path(&manifest).ends_with(
"flat/refs/sha256_aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa.json"
));
assert!(cache.flat_blob_path(&artifact).ends_with(
"flat/blobs/sha256_cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc.raw"
));
assert!(cache.flat_lock_path(&derivation).ends_with(
"flat/locks/sha256_bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb.lock"
));
let blob_path = cache.flat_blob_path(&artifact);
let blob = std::fs::File::create(&blob_path).unwrap();
blob.set_len(4096).unwrap();
let reference = FlatRootfsRef {
schema: 1,
manifest_digest: manifest.to_string(),
derivation_digest: derivation.to_string(),
artifact_digest: artifact.to_string(),
materializer_abi: 1,
uuid: "00".repeat(16),
virtual_size_bytes: 4096,
inode_count: 2,
content_bytes: 7,
};
cache.write_flat_ref(&manifest, &reference).unwrap();
assert_eq!(cache.read_flat_ref(&manifest).unwrap(), Some(reference));
}
#[test]
fn aligned_garbage_is_not_a_valid_erofs_cache_artifact() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("corrupt.erofs");
let file = std::fs::File::create(&path).unwrap();
file.set_len(EROFS_ALIGNMENT_BYTES).unwrap();
assert!(!is_valid_erofs_artifact(&path));
}
#[test]
fn flat_ref_must_name_the_manifest_that_indexes_it() {
let directory = tempfile::tempdir().unwrap();
let cache = GlobalCache::new(directory.path()).unwrap();
let indexed_manifest = digest('a');
let wrong_manifest = digest('b');
let artifact = digest('c');
std::fs::write(cache.flat_blob_path(&artifact), [0u8; 8]).unwrap();
let reference = FlatRootfsRef {
schema: 1,
manifest_digest: wrong_manifest.to_string(),
derivation_digest: digest('d').to_string(),
artifact_digest: artifact.to_string(),
materializer_abi: 1,
uuid: "00".repeat(16),
virtual_size_bytes: 8,
inode_count: 2,
content_bytes: 0,
};
cache.write_flat_ref(&indexed_manifest, &reference).unwrap();
assert_eq!(cache.read_flat_ref(&indexed_manifest).unwrap(), None);
}
#[test]
fn publishing_replaces_same_size_blob_with_wrong_content() {
let directory = tempfile::tempdir().unwrap();
let cache = GlobalCache::new(directory.path()).unwrap();
let candidate = directory.path().join("candidate.raw");
let expected_bytes = b"good";
std::fs::write(&candidate, expected_bytes).unwrap();
let expected: Digest = format!("sha256:{}", hex::encode(Sha256::digest(expected_bytes)))
.parse()
.unwrap();
let destination = cache.flat_blob_path(&expected);
std::fs::write(&destination, b"evil").unwrap();
cache
.publish_flat_blob(&candidate, &expected, expected_bytes.len() as u64)
.unwrap();
assert_eq!(std::fs::read(destination).unwrap(), expected_bytes);
}
}