use super::index::SqliteIndexStore;
use super::RefUpdate;
use crate::artifact::{
media_types::{self, RootPayloadVersion},
sha256_digest, stable_json_bytes, ImageRef,
};
use anyhow::{ensure, Context, Result};
use oci_spec::image::{Descriptor, DescriptorBuilder, Digest, ImageManifest, MediaType};
use std::collections::{BTreeSet, HashMap};
use std::ops::Deref;
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::sync::OnceLock;
mod blob;
mod gc;
mod import;
use blob::{BlobRecord, DeleteBlobOutcome, FileBlobStore};
pub use gc::{
GcBlob, GcDeleteReport, GcInvalidManifest, GcMissingBlob, GcOptions, GcReferenceKind, GcReport,
GcRoot,
};
pub use import::{ArchiveInspectView, LegacyImportReport, OciDirImport, OciDirRef};
static DEFAULT_LOCAL_REGISTRY: OnceLock<LocalRegistry> = OnceLock::new();
const EXPERIMENT_CHECKPOINT_REPOSITORY: &str = "checkpoint";
const FILE_BLOB_STORE_DIR_NAME: &str = "blobs";
#[derive(Debug, Clone)]
pub struct StoredDescriptor<'reg> {
registry: &'reg LocalRegistry,
descriptor: Descriptor,
}
impl StoredDescriptor<'_> {
pub fn ensure_media_type(&self, expected: &MediaType) -> Result<()> {
let actual = self.media_type();
ensure!(
actual == expected,
"Expected media type '{expected}', got '{actual}'"
);
Ok(())
}
pub(crate) fn is_stored_in(&self, registry: &LocalRegistry) -> bool {
std::ptr::eq(self.registry, registry)
}
pub(crate) fn registry(&self) -> &LocalRegistry {
self.registry
}
fn into_inner(self) -> Descriptor {
self.descriptor
}
}
impl Deref for StoredDescriptor<'_> {
type Target = Descriptor;
fn deref(&self) -> &Self::Target {
&self.descriptor
}
}
impl From<StoredDescriptor<'_>> for Descriptor {
fn from(value: StoredDescriptor<'_>) -> Self {
value.into_inner()
}
}
#[derive(Debug, Clone)]
pub(crate) struct SealedArtifact<'reg>(StoredDescriptor<'reg>);
impl<'reg> Deref for SealedArtifact<'reg> {
type Target = StoredDescriptor<'reg>;
fn deref(&self) -> &Self::Target {
&self.0
}
}
impl SealedArtifact<'_> {
fn is_stored_in(&self, registry: &LocalRegistry) -> bool {
self.0.is_stored_in(registry)
}
}
#[derive(Debug, Clone)]
pub(crate) struct UnsealedArtifact<'reg> {
artifact_type: MediaType,
config: StoredDescriptor<'reg>,
layers: Vec<StoredDescriptor<'reg>>,
subject: Option<Descriptor>,
annotations: HashMap<String, String>,
}
impl<'reg> UnsealedArtifact<'reg> {
pub(crate) fn new(
artifact_type: MediaType,
config: StoredDescriptor<'reg>,
layers: Vec<StoredDescriptor<'reg>>,
subject: Option<Descriptor>,
annotations: HashMap<String, String>,
) -> Self {
Self {
artifact_type,
config,
layers,
subject,
annotations,
}
}
pub(crate) fn into_oci_image_manifest(self) -> Result<ImageManifest> {
let config: Descriptor = self.config.into();
let mut builder = oci_spec::image::ImageManifestBuilder::default()
.schema_version(2u32)
.artifact_type(self.artifact_type)
.config(config)
.layers(self.layers.into_iter().map(Into::into).collect::<Vec<_>>());
if let Some(subject) = self.subject {
builder = builder.subject(subject);
}
if !self.annotations.is_empty() {
builder = builder.annotations(self.annotations);
}
builder
.build()
.context("Failed to build OCI image manifest")
}
fn ensure_stored_in(&self, registry: &LocalRegistry) -> Result<()> {
ensure!(
self.config.is_stored_in(registry),
"Artifact config descriptor belongs to a different Local Registry"
);
ensure!(
self.layers
.iter()
.all(|descriptor| descriptor.is_stored_in(registry)),
"Artifact layer descriptor belongs to a different Local Registry"
);
Ok(())
}
}
#[derive(Debug)]
pub struct LocalRegistry {
root: PathBuf,
index: SqliteIndexStore,
blobs: FileBlobStore,
}
#[derive(Debug)]
pub struct TempLocalRegistry {
registry: LocalRegistry,
tempdir: tempfile::TempDir,
}
impl TempLocalRegistry {
pub fn new() -> Result<Self> {
let tempdir = tempfile::tempdir().context("Failed to create temporary Local Registry")?;
let registry = LocalRegistry::open(tempdir.path())?;
Ok(Self { registry, tempdir })
}
pub fn registry(&self) -> &LocalRegistry {
&self.registry
}
pub fn path(&self) -> &Path {
self.tempdir.path()
}
}
impl LocalRegistry {
pub fn open(root: impl Into<PathBuf>) -> Result<Self> {
let root = root.into();
let index = SqliteIndexStore::open_in_registry_root(&root)?;
let blobs = FileBlobStore::new(root.join(FILE_BLOB_STORE_DIR_NAME))?;
Ok(Self { root, index, blobs })
}
pub fn open_default() -> Result<Self> {
Self::open(crate::artifact::get_local_registry_root())
}
pub fn shared_default() -> Result<&'static Self> {
if let Some(registry) = DEFAULT_LOCAL_REGISTRY.get() {
return Ok(registry);
}
let registry = Self::open_default()?;
let _ = DEFAULT_LOCAL_REGISTRY.set(registry);
Ok(DEFAULT_LOCAL_REGISTRY
.get()
.expect("default Local Registry was initialized"))
}
pub fn root(&self) -> &Path {
&self.root
}
pub fn get_blob(&self, descriptor: &StoredDescriptor<'_>) -> Result<Vec<u8>> {
ensure!(
descriptor.is_stored_in(self),
"Descriptor {} is not stored in this Local Registry",
descriptor.digest()
);
let bytes = self.read_blob(descriptor.digest())?;
ensure!(
bytes.len() as u64 == descriptor.size(),
"Descriptor size mismatch for {}: descriptor={}, actual={}",
descriptor.digest(),
descriptor.size(),
bytes.len()
);
Ok(bytes)
}
pub fn get_instance_layer(&self, descriptor: &StoredDescriptor<'_>) -> Result<crate::Instance> {
let payload_version = media_types::instance_payload_version(descriptor.media_type())?;
let bytes = self.get_blob(descriptor)?;
let annotations = descriptor
.annotations()
.as_ref()
.cloned()
.unwrap_or_default();
let mut instance = match payload_version {
RootPayloadVersion::V1 => crate::Instance::from_v1_bytes(&bytes)?,
RootPayloadVersion::V2 => crate::Instance::from_v2_bytes(&bytes)?,
};
crate::FlatAnnotations::merge_annotations(&mut instance, &annotations);
Ok(instance)
}
pub fn get_parametric_instance_layer(
&self,
descriptor: &StoredDescriptor<'_>,
) -> Result<crate::ParametricInstance> {
let payload_version =
media_types::parametric_instance_payload_version(descriptor.media_type())?;
let bytes = self.get_blob(descriptor)?;
let annotations = descriptor
.annotations()
.as_ref()
.cloned()
.unwrap_or_default();
let mut instance = match payload_version {
RootPayloadVersion::V1 => crate::ParametricInstance::from_v1_bytes(&bytes)?,
RootPayloadVersion::V2 => crate::ParametricInstance::from_v2_bytes(&bytes)?,
};
crate::FlatAnnotations::merge_annotations(&mut instance, &annotations);
Ok(instance)
}
pub fn get_solution_layer(&self, descriptor: &StoredDescriptor<'_>) -> Result<crate::Solution> {
let payload_version = media_types::solution_payload_version(descriptor.media_type())?;
let bytes = self.get_blob(descriptor)?;
let annotations = descriptor
.annotations()
.as_ref()
.cloned()
.unwrap_or_default();
let mut solution = match payload_version {
RootPayloadVersion::V1 => crate::Solution::from_v1_bytes(&bytes)?,
RootPayloadVersion::V2 => crate::Solution::from_v2_bytes(&bytes)?,
};
crate::FlatAnnotations::merge_annotations(&mut solution, &annotations);
Ok(solution)
}
pub fn get_sample_set_layer(
&self,
descriptor: &StoredDescriptor<'_>,
) -> Result<crate::SampleSet> {
let payload_version = media_types::sample_set_payload_version(descriptor.media_type())?;
let bytes = self.get_blob(descriptor)?;
let annotations = descriptor
.annotations()
.as_ref()
.cloned()
.unwrap_or_default();
let mut sample_set = match payload_version {
RootPayloadVersion::V1 => crate::SampleSet::from_v1_bytes(&bytes)?,
RootPayloadVersion::V2 => crate::SampleSet::from_v2_bytes(&bytes)?,
};
crate::FlatAnnotations::merge_annotations(&mut sample_set, &annotations);
Ok(sample_set)
}
pub fn store_instance_layer(&self, instance: &crate::Instance) -> Result<StoredDescriptor<'_>> {
self.store_layer_blob(
media_types::v2_instance(),
&instance.to_v2_bytes(),
crate::FlatAnnotations::flat_annotations(instance),
)
}
pub fn store_parametric_instance_layer(
&self,
instance: &crate::ParametricInstance,
) -> Result<StoredDescriptor<'_>> {
self.store_layer_blob(
media_types::v2_parametric_instance(),
&instance.to_v2_bytes(),
crate::FlatAnnotations::flat_annotations(instance),
)
}
pub fn store_solution_layer(&self, solution: &crate::Solution) -> Result<StoredDescriptor<'_>> {
self.store_layer_blob(
media_types::v2_solution(),
&solution.to_v2_bytes(),
crate::FlatAnnotations::flat_annotations(solution),
)
}
pub fn store_sample_set_layer(
&self,
sample_set: &crate::SampleSet,
) -> Result<StoredDescriptor<'_>> {
self.store_layer_blob(
media_types::v2_sample_set(),
&sample_set.to_v2_bytes(),
crate::FlatAnnotations::flat_annotations(sample_set),
)
}
pub fn resolve_image_name(&self, image_name: &ImageRef) -> Result<Option<Digest>> {
self.index.resolve_image_name(image_name)
}
pub(crate) fn registry_id(&self) -> Result<String> {
self.index.registry_id()
}
pub fn synthesize_anonymous_image_name(&self) -> Result<ImageRef> {
let registry_id = self.index.registry_id()?;
crate::artifact::anonymous_artifact_image_name(®istry_id)
}
pub fn synthesize_anonymous_experiment_image_name(&self) -> Result<ImageRef> {
let registry_id = self.index.registry_id()?;
crate::artifact::anonymous_local_image_name(®istry_id, "experiment")
.with_context(|| "Failed to synthesise anonymous experiment image name")
}
pub(crate) fn experiment_checkpoint_image_name(
&self,
requested_image_name: &ImageRef,
) -> Result<ImageRef> {
let registry_id = self.index.registry_id()?;
let repository_key = crate::artifact::anonymous_local_repository_key(
®istry_id,
EXPERIMENT_CHECKPOINT_REPOSITORY,
)?;
let digest = sha256_digest(requested_image_name.to_string().as_bytes());
let tag = digest
.strip_prefix("sha256:")
.expect("sha256_digest returns a sha256-prefixed digest");
ImageRef::parse(&format!("{repository_key}:{tag}")).with_context(|| {
format!("Failed to derive experiment checkpoint image name for {requested_image_name}")
})
}
pub fn list_anonymous_artifact_refs(
&self,
) -> Result<Vec<crate::artifact::local_registry::RefRecord>> {
let all = self.index.list_refs(None)?;
Ok(all
.into_iter()
.filter(|r| {
crate::artifact::is_anonymous_artifact_ref_name(&r.name)
&& crate::artifact::is_anonymous_artifact_tag(&r.reference)
})
.collect())
}
pub fn prune_anonymous_artifact_refs(
&self,
) -> Result<Vec<crate::artifact::local_registry::RefRecord>> {
let refs = self.list_anonymous_artifact_refs()?;
for r in &refs {
self.index.delete_ref(&r.name, &r.reference)?;
}
Ok(refs)
}
pub fn list_image_refs(&self) -> Result<Vec<ImageRef>> {
self.index
.list_refs(None)?
.into_iter()
.map(|r| ImageRef::from_repository_and_reference(&r.name, &r.reference))
.collect()
}
pub(crate) fn seal_artifact<'reg>(
&'reg self,
artifact: UnsealedArtifact<'reg>,
) -> Result<SealedArtifact<'reg>> {
artifact.ensure_stored_in(self)?;
let manifest = artifact.into_oci_image_manifest()?;
Self::validate_manifest(&manifest)?;
let manifest_bytes = stable_json_bytes(&manifest)?;
let manifest_descriptor = Self::build_manifest_descriptor(&manifest_bytes)?;
let stored_manifest = self.store_blob(manifest_descriptor, &manifest_bytes)?;
Ok(SealedArtifact(stored_manifest))
}
pub(crate) fn publish_manifest_ref(
&self,
image_name: &ImageRef,
sealed_artifact: &SealedArtifact<'_>,
) -> Result<RefUpdate> {
ensure!(
sealed_artifact.is_stored_in(self),
"Sealed artifact descriptor belongs to a different Local Registry"
);
self.index.publish_image_ref(image_name, &sealed_artifact.0)
}
pub(crate) fn publish_stored_manifest_ref(
&self,
image_name: &ImageRef,
manifest: &StoredDescriptor<'_>,
) -> Result<RefUpdate> {
ensure!(
manifest.is_stored_in(self),
"Manifest descriptor belongs to a different Local Registry"
);
self.index.publish_image_ref(image_name, manifest)
}
pub(crate) fn replace_manifest_ref(
&self,
image_name: &ImageRef,
sealed_artifact: &SealedArtifact<'_>,
) -> Result<RefUpdate> {
ensure!(
sealed_artifact.is_stored_in(self),
"Sealed artifact descriptor belongs to a different Local Registry"
);
self.index.replace_image_ref(image_name, &sealed_artifact.0)
}
pub(crate) fn delete_manifest_ref(&self, image_name: &ImageRef) -> Result<bool> {
self.index
.delete_ref(&image_name.repository_key(), image_name.reference())
}
fn store_blob_bytes(&self, bytes: &[u8]) -> Result<Digest> {
self.blobs.put_bytes(bytes)
}
pub(crate) fn read_blob(&self, digest: &Digest) -> Result<Vec<u8>> {
self.blobs.read_bytes(digest)
}
pub(crate) fn contains_blob(&self, digest: &Digest) -> Result<bool> {
self.blobs.exists(digest)
}
pub(crate) fn blob_size(&self, digest: &Digest) -> Result<u64> {
self.blobs.size(digest)
}
pub(crate) fn touch_blob(&self, digest: &Digest) -> Result<()> {
self.blobs.touch_blob(digest)
}
fn list_blob_records(&self) -> Result<Vec<BlobRecord>> {
self.blobs.list_blobs()
}
fn delete_blob_if_older_than(
&self,
digest: &Digest,
cutoff: std::time::SystemTime,
) -> Result<DeleteBlobOutcome> {
self.blobs.delete_blob_if_older_than(digest, cutoff)
}
pub(crate) fn stored_manifest_descriptor(
&self,
manifest_digest: &Digest,
) -> Result<StoredDescriptor<'_>> {
let size = self.blob_size(manifest_digest)?;
let descriptor = DescriptorBuilder::default()
.media_type(MediaType::ImageManifest)
.digest(manifest_digest.clone())
.size(size)
.build()
.context("Failed to build manifest descriptor")?;
self.stored_descriptor(descriptor)
}
pub(crate) fn touch_manifest_closure(
&self,
manifest_digest: &Digest,
visited: &mut BTreeSet<String>,
) -> Result<()> {
if !visited.insert(manifest_digest.as_ref().to_string()) {
return Ok(());
}
self.touch_blob(manifest_digest)?;
let bytes = self
.read_blob(manifest_digest)
.with_context(|| format!("Failed to read manifest blob {manifest_digest}"))?;
let manifest: ImageManifest = serde_json::from_slice(&bytes)
.with_context(|| format!("Failed to parse OCI image manifest {manifest_digest}"))?;
self.touch_descriptor_blob(manifest.config())?;
for layer in manifest.layers() {
self.touch_descriptor_blob(layer)?;
}
if let Some(subject) = manifest.subject() {
let subject = self.stored_descriptor(subject.clone())?;
self.touch_manifest_closure(subject.digest(), visited)?;
}
Ok(())
}
fn touch_descriptor_blob(&self, descriptor: &Descriptor) -> Result<()> {
let descriptor = self.stored_descriptor(descriptor.clone())?;
self.touch_blob(descriptor.digest())
}
fn validate_manifest(manifest: &ImageManifest) -> Result<()> {
let artifact_type = manifest
.artifact_type()
.as_ref()
.context("Manifest does not carry the OMMX `artifactType` field")?;
ensure!(
artifact_type == &MediaType::Other(media_types::V1_ARTIFACT_MEDIA_TYPE.to_string()),
"Manifest `artifactType` must be `{}`, got `{}`",
media_types::V1_ARTIFACT_MEDIA_TYPE,
artifact_type,
);
Ok(())
}
fn build_manifest_descriptor(manifest_bytes: &[u8]) -> Result<Descriptor> {
DescriptorBuilder::default()
.media_type(MediaType::ImageManifest)
.digest(
Digest::from_str(&sha256_digest(manifest_bytes))
.context("Failed to parse manifest digest")?,
)
.size(manifest_bytes.len() as u64)
.build()
.context("Failed to build manifest descriptor")
}
pub(crate) fn store_layer_blob(
&self,
media_type: MediaType,
bytes: &[u8],
annotations: HashMap<String, String>,
) -> Result<StoredDescriptor<'_>> {
let digest =
Digest::from_str(&sha256_digest(bytes)).context("Failed to parse layer blob digest")?;
let descriptor = DescriptorBuilder::default()
.media_type(media_type)
.digest(digest)
.size(bytes.len() as u64)
.annotations(annotations)
.build()
.context("Failed to build layer descriptor")?;
self.store_blob(descriptor, bytes)
}
pub(crate) fn store_json_layer_blob(
&self,
media_type: MediaType,
value: &impl serde::Serialize,
annotations: HashMap<String, String>,
) -> Result<StoredDescriptor<'_>> {
let bytes = serde_json::to_vec(value).context("Failed to encode JSON layer")?;
self.store_layer_blob(media_type, &bytes, annotations)
}
pub(crate) fn store_json_blob(
&self,
media_type: MediaType,
value: &impl serde::Serialize,
) -> Result<StoredDescriptor<'_>> {
let bytes = serde_json::to_vec(value).context("Failed to encode JSON blob")?;
let digest =
Digest::from_str(&sha256_digest(&bytes)).context("Failed to parse JSON blob digest")?;
let descriptor = DescriptorBuilder::default()
.media_type(media_type)
.digest(digest)
.size(bytes.len() as u64)
.build()
.context("Failed to build JSON blob descriptor")?;
self.store_blob(descriptor, &bytes)
}
pub(crate) fn store_blob(
&self,
descriptor: Descriptor,
bytes: &[u8],
) -> Result<StoredDescriptor<'_>> {
let digest = self.store_blob_bytes(bytes)?;
ensure!(
&digest == descriptor.digest(),
"Descriptor digest mismatch: descriptor={}, actual={}",
descriptor.digest(),
digest
);
ensure!(
bytes.len() as u64 == descriptor.size(),
"Descriptor size mismatch for {}: descriptor={}, actual={}",
descriptor.digest(),
descriptor.size(),
bytes.len()
);
Ok(StoredDescriptor {
registry: self,
descriptor,
})
}
pub(crate) fn stored_descriptor(&self, descriptor: Descriptor) -> Result<StoredDescriptor<'_>> {
let size = self.blob_size(descriptor.digest())?;
ensure!(
size == descriptor.size(),
"Descriptor size mismatch for {}: descriptor={}, actual={}",
descriptor.digest(),
descriptor.size(),
size
);
Ok(StoredDescriptor {
registry: self,
descriptor,
})
}
}