use super::{
manifest::{blob_descriptors_by_digest, read_blob_by_descriptor_async},
remote_transport::{blob_transfer_concurrency, bounded_map, RemoteTransport},
LocalArtifact, LocalManifest,
};
use anyhow::Context;
use oci_client::RegistryOperation;
use oci_spec::image::Descriptor;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum PushBlobOutcome {
Skipped,
Transferred,
}
impl LocalArtifact<'_> {
pub fn push(&self) -> crate::Result<()> {
let manifest = self.get_manifest()?.clone();
let blob_descriptors = collect_blob_descriptors(&manifest);
let descriptor_count = blob_descriptors.len();
let unique_blob_descriptors = blob_descriptors_by_digest(&blob_descriptors)?;
let deduplicated = descriptor_count - unique_blob_descriptors.len();
let transfer_concurrency = blob_transfer_concurrency()?;
let transport = RemoteTransport::new(self.image_name())?;
transport.auth_for(self.image_name(), RegistryOperation::Pull)?;
transport.auth(self.image_name())?;
let outcomes = transport.block_on(async {
bounded_map(
unique_blob_descriptors.into_values(),
transfer_concurrency,
|descriptor| {
let transport = &transport;
push_descriptor_if_missing(
descriptor,
move |digest| async move {
transport
.blob_exists_async(self.image_name(), &digest)
.await
.with_context(|| format!("Failed while checking blob {digest}"))
},
|descriptor| async move {
read_blob_by_descriptor_async(self, &descriptor).await
},
move |digest, bytes| async move {
transport
.push_blob_async(self.image_name(), &digest, bytes)
.await
.with_context(|| format!("Failed while uploading blob {digest}"))
},
)
},
)
.await
})?;
let checked = outcomes.len();
let skipped = outcomes
.iter()
.filter(|outcome| **outcome == PushBlobOutcome::Skipped)
.count();
let transferred = checked - skipped;
tracing::info!(
checked,
skipped,
transferred,
deduplicated,
concurrency = transfer_concurrency.get(),
"Completed remote blob transfers"
);
let manifest_bytes = self.read_blob_by_digest(self.manifest_digest())?;
let content_type = manifest.media_type();
tracing::info!(
"Publishing manifest {} ({}, {} bytes) to {}",
self.manifest_digest(),
content_type,
manifest_bytes.len(),
self.image_name(),
);
transport.push_manifest_bytes(self.image_name(), manifest_bytes, content_type)?;
Ok(())
}
}
async fn push_descriptor_if_missing<Exists, ExistsFuture, Read, ReadFuture, Push, PushFuture>(
descriptor: &Descriptor,
blob_exists: Exists,
read_blob: Read,
push_blob: Push,
) -> crate::Result<PushBlobOutcome>
where
Exists: FnOnce(String) -> ExistsFuture,
ExistsFuture: std::future::Future<Output = crate::Result<bool>>,
Read: FnOnce(Descriptor) -> ReadFuture,
ReadFuture: std::future::Future<Output = crate::Result<Vec<u8>>>,
Push: FnOnce(String, Vec<u8>) -> PushFuture,
PushFuture: std::future::Future<Output = crate::Result<()>>,
{
let digest = descriptor.digest().to_string();
if blob_exists(digest.clone()).await? {
tracing::debug!(%digest, "Skipping blob already present in remote");
return Ok(PushBlobOutcome::Skipped);
}
let bytes = read_blob(descriptor.clone()).await?;
tracing::debug!(size = bytes.len(), %digest, "Pushing blob");
push_blob(digest, bytes).await?;
Ok(PushBlobOutcome::Transferred)
}
fn collect_blob_descriptors(manifest: &LocalManifest) -> Vec<Descriptor> {
let layers = manifest.layers();
let mut out = Vec::with_capacity(1 + layers.len());
out.push(manifest.config());
out.extend(layers);
out
}
#[cfg(test)]
mod tests {
use super::*;
use oci_spec::image::{DescriptorBuilder, Digest, MediaType};
use std::{
cell::{Cell, RefCell},
rc::Rc,
str::FromStr,
};
fn descriptor_for(bytes: &[u8]) -> Descriptor {
DescriptorBuilder::default()
.media_type(MediaType::Other("application/octet-stream".to_string()))
.digest(Digest::from_str(&crate::artifact::sha256_digest(bytes)).unwrap())
.size(bytes.len() as u64)
.build()
.unwrap()
}
#[test]
fn already_present_descriptors_are_not_uploaded() {
let present = descriptor_for(b"present");
let read_count = Cell::new(0);
let push_count = Cell::new(0);
let runtime = tokio::runtime::Builder::new_current_thread()
.build()
.unwrap();
let push_count_ref = &push_count;
let outcome = runtime
.block_on(push_descriptor_if_missing(
&present,
|_| async { Ok(true) },
|_| async {
read_count.set(read_count.get() + 1);
Ok(Vec::new())
},
move |_, _| async move {
push_count_ref.set(push_count_ref.get() + 1);
Ok(())
},
))
.unwrap();
assert_eq!(outcome, PushBlobOutcome::Skipped);
assert_eq!(read_count.get(), 0);
assert_eq!(push_count.get(), 0);
}
#[test]
fn missing_descriptor_is_read_and_uploaded_in_one_operation() {
let missing = descriptor_for(b"missing");
let events = Rc::new(RefCell::new(Vec::new()));
let runtime = tokio::runtime::Builder::new_current_thread()
.build()
.unwrap();
let check_events = Rc::clone(&events);
let read_events = Rc::clone(&events);
let push_events = Rc::clone(&events);
let outcome = runtime
.block_on(push_descriptor_if_missing(
&missing,
move |_| async move {
check_events.borrow_mut().push("check");
Ok(false)
},
move |_| async move {
read_events.borrow_mut().push("read");
Ok(b"missing".to_vec())
},
move |_, bytes| async move {
assert_eq!(bytes, b"missing");
push_events.borrow_mut().push("push");
Ok(())
},
))
.unwrap();
assert_eq!(outcome, PushBlobOutcome::Transferred);
assert_eq!(*events.borrow(), ["check", "read", "push"]);
}
}