use std::collections::HashMap;
use bollard::models::ImageDeleteResponseItem;
use bollard::models::ImageSummary;
use bollard::query_parameters::CreateImageOptions;
use bollard::query_parameters::ListImagesOptions;
use bollard::query_parameters::RemoveImageOptions;
use crankshaft_events::Event;
use crankshaft_events::TaskId;
use crankshaft_events::send_event;
use futures::stream::FuturesUnordered;
use tokio::sync::broadcast;
use tokio_stream::StreamExt as _;
use tokio_util::sync::CancellationToken;
use tracing::Level;
use tracing::debug;
use tracing::enabled;
use tracing::trace;
use crate::Docker;
use crate::Error;
use crate::Result;
pub(crate) async fn list_images(docker: &Docker) -> Result<Vec<ImageSummary>> {
debug!("listing images");
let images = docker
.inner()
.list_images(Some(ListImagesOptions {
all: true,
..Default::default()
}))
.await
.map_err(Error::Docker)?;
debug!("found {} images", images.len());
if enabled!(Level::TRACE) {
for image in &images {
trace!(
" image: {} (tags: {})",
image.id,
image
.repo_tags
.iter()
.map(|v| format!("`{v}`"))
.collect::<Vec<_>>()
.join(", ")
);
for tag in &image.repo_tags {
trace!(" {}", tag);
}
}
}
Ok(images)
}
pub(crate) async fn ensure_image(
docker: &Docker,
image: impl Into<String>,
token: CancellationToken,
events_ctx: Option<(broadcast::Sender<Event>, TaskId)>,
) -> Result<Option<()>> {
let image = image.into();
debug!("ensuring image `{image}` exists locally");
let mut filters = HashMap::new();
filters.insert(String::from("reference"), vec![image.clone()]);
let results = docker
.inner()
.list_images(Some(ListImagesOptions {
filters: Some(filters),
..Default::default()
}))
.await
.map_err(Error::Docker)?;
if !results.is_empty() {
debug!("image `{image}` exists locally");
if enabled!(Level::TRACE) {
trace!(
"image SHA = {}",
results.first().unwrap().id.trim_start_matches("sha256:")
);
}
return Ok(Some(()));
}
let events = events_ctx.as_ref().map(|(sender, _)| sender.clone());
let task = events_ctx.as_ref().map(|(_, task_id)| *task_id);
debug!("image `{image}` does not exist locally; attempting to pull from remote");
send_event!(
events,
Event::ImagePullStarted {
id: task.unwrap(),
name: image.clone()
}
);
let mut stream = docker.inner().create_image(
Some(CreateImageOptions {
tag: Some(if image.contains(':') {
String::from("")
} else {
String::from("latest")
}),
from_image: Some(image.clone()),
..Default::default()
}),
None,
None,
);
loop {
tokio::select! {
biased;
_ = token.cancelled() => {
debug!("image pull cancelled");
return Ok(None);
}
item = stream.next() => {
let Some(result) = item else {
break;
};
let update = match result {
Ok(update) => update,
Err(e) => {
send_event!(
events,
Event::ImagePullFailed {
id: task.unwrap(),
name: image.clone(),
message: e.to_string()
}
);
return Err(Error::Docker(e));
}
};
trace!(
"pull update: {}",
[
update.id.map(|id| format!("id: {id}")),
update.error_detail.and_then(|err| err.message).map(|err| format!("error: {err}")),
update.status.map(|status| format!("status: {status}")),
update.progress_detail.map(|progress|
format!(
"progress: {}/{}",
progress
.current
.map(|v| v.to_string())
.unwrap_or(String::from("?")),
progress
.total
.map(|v| v.to_string())
.unwrap_or(String::from("?"))
))
]
.into_iter()
.flatten()
.collect::<Vec<_>>()
.join("; ")
)
}
}
}
send_event!(
events,
Event::ImagePullFinished {
id: task.unwrap(),
name: image
}
);
Ok(Some(()))
}
pub(crate) async fn remove_image<T: AsRef<str>, U: AsRef<str>>(
docker: &Docker,
name: T,
tag: U,
) -> Result<impl IntoIterator<Item = ImageDeleteResponseItem> + use<T, U>> {
let name = name.as_ref();
let tag = tag.as_ref();
debug!("removing image: {name} ({tag})");
let images = docker
.inner()
.remove_image(name, None::<RemoveImageOptions>, None)
.await
.map_err(Error::Docker)?;
if enabled!(Level::TRACE) {
for image in &images {
if let Some(untagged) = &image.untagged {
trace!(" untagged image: {untagged}");
}
if let Some(deleted) = &image.deleted {
trace!(" deleted image: {deleted}");
}
}
}
Ok(images)
}
pub(crate) async fn remove_all_images(docker: &Docker) -> Result<Vec<ImageDeleteResponseItem>> {
debug!("removing all images");
let mut results = Vec::new();
for image in docker.list_images().await? {
let mut futures = FuturesUnordered::new();
for tag in image.repo_tags {
futures.push(docker.remove_image(&image.id, tag));
}
while let Some(result) = futures.next().await {
results.extend(result?);
}
}
if !results.is_empty() {
debug!("removed {} images in total", results.len());
}
Ok(results)
}