use std::collections::HashMap;
use std::io;
use std::path::Path;
use std::sync::Arc;
use crate::image::{
oci::{
digest::Digest,
layout::OCIImageLayout,
spec_v1::{Descriptor, Image as OCIImage, Index, Manifest},
},
transports,
types::ImageSource,
};
use tokio::{io::BufReader, sync::Semaphore};
pub async fn pull_container_image<P>(
reference: &str,
to_path: P,
force: bool,
clean_on_err: bool,
) -> std::io::Result<OCIImageLayout>
where
P: AsRef<Path> + std::fmt::Debug,
{
log::info!("Pulling the image: {}", reference);
let image_ref = transports::parse_image_name(reference)?;
let docker_ref = image_ref.docker_reference();
if docker_ref.is_none() {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
format!("Invalid Image Name {}", reference),
));
}
let name = docker_ref.as_ref().unwrap().name();
let tag = docker_ref.as_ref().unwrap().tag();
log::debug!(
"Creating OCI Image Layout for Image: {}, {}, {:?}",
&name,
&tag,
to_path
);
let mut img_layout = OCIImageLayout::new(&name, Some(&tag), to_path);
if img_layout.image_fs_path().exists() {
if !force {
let errstr = format!("Local FS path for the image with name: {}, tag: {} exists. Please specify `--force` to overwrite.", name, tag);
log::error!("{}", errstr);
return Err(io::Error::new(io::ErrorKind::InvalidInput, errstr));
} else {
log::warn!("Local Image Layout exists, User requested 'force'. Deleting...");
img_layout.delete_fs_path().await?;
}
}
img_layout.create_fs_path().await?;
log::debug!("Performing Image Pull.");
let result = match perform_image_pull(&mut img_layout, reference).await {
Ok(_) => Ok(img_layout),
Err(e) => {
eprintln!("Error : {}", e);
if clean_on_err {
img_layout.delete_fs_path().await?;
}
Err(e)
}
};
result
}
async fn perform_image_pull(
img_layout: &mut OCIImageLayout,
image_name: &str,
) -> std::io::Result<()> {
let image_ref = transports::parse_image_name(image_name)?;
let mut img = image_ref.new_image()?;
log::trace!("Getting Manifest for the Image.");
let manifest = img.resolved_manifest().await?;
log::trace!("Writing Manifest Blob.");
let digest = Digest::from_bytes(&manifest.manifest);
let mut reader = BufReader::new(&*manifest.manifest);
img_layout.write_blob_file(&digest, &mut reader).await?;
let mut annotations = HashMap::new();
let _ = annotations.insert(
"org.opencontainers.image.ref.name".to_string(),
img_layout.tag().as_ref().unwrap().clone(),
);
let manifest_descriptor = Descriptor {
mediatype: Some(manifest.mime_type.to_string()),
digest: digest,
size: manifest.manifest.len() as i64,
urls: None,
platform: None,
annotations: Some(annotations),
};
log::trace!("Updating Image Layout 'Index', with new manifest.");
img_layout.update_index(Index {
version: 2,
manifests: vec![manifest_descriptor],
annotations: None,
});
log::trace!("Getting Image Config.");
let manifest_obj: Manifest = serde_json::from_slice(&manifest.manifest)?;
log::trace!("Saving Image Config.");
let config = img.config_blob().await?;
let mut reader = BufReader::new(&*config);
img_layout
.write_blob_file(&manifest_obj.config.digest, &mut reader)
.await?;
let image_obj: OCIImage = serde_json::from_slice(&config)?;
log::debug!("Getting Image Layers!");
let max_parallel_dloads = 3;
let mut layer_handles = vec![];
let semaphore = Arc::new(Semaphore::new(max_parallel_dloads));
for (layer, unzipped_digest) in manifest_obj.layers.iter().zip(image_obj.rootfs.diff_ids) {
let layer_digest = layer.digest.clone();
let img_layout = img_layout.clone();
let img_source = image_ref.new_image_source()?;
let permit = semaphore.clone().acquire_owned().await;
let handle = tokio::spawn(async move {
do_download_image_layer(layer_digest, unzipped_digest, img_layout, img_source).await?;
drop(permit);
Ok::<(), std::io::Error>(())
});
layer_handles.push(handle);
}
for h in layer_handles {
let _ = h.await?;
}
log::debug!("Writing 'index.json'.");
img_layout.write_index_json().await?;
log::debug!("Writing 'img-layout'.");
img_layout.write_image_layout().await?;
log::info!("Image downloaded and saved successfully!");
Ok(())
}
async fn do_download_image_layer<'a>(
layer_digest: Digest,
unzipped_digest: Digest,
img_layout: OCIImageLayout,
img_source: Box<dyn ImageSource + Send + Sync>,
) -> io::Result<()> {
log::info!("Getting Image Layer: {}", layer_digest);
let layer_reader = img_source.get_blob(&layer_digest).await?;
log::trace!("Layer downloaded, Verifying the RootFS Layer.");
let reader = BufReader::new(layer_reader);
let mut gzip_decoder = async_compression::tokio::bufread::GzipDecoder::new(reader);
let unzipped_verify = unzipped_digest.verify(&mut gzip_decoder).await;
if unzipped_verify {
log::trace!("Image Layer {} verified. Saving Image Layer.", layer_digest);
let layer_reader = img_source.get_blob(&layer_digest).await?;
let mut reader = BufReader::new(layer_reader);
&img_layout
.write_blob_file(&layer_digest, &mut reader)
.await?;
} else {
log::error!(
"Checksum does not match for: {} after uncompressing.",
&layer_digest
);
}
Ok(())
}