use std::collections::HashMap;
use async_trait::async_trait;
use bollard::Docker;
use bollard::query_parameters::{
CreateImageOptions, ListImagesOptions, PushImageOptions, RemoveImageOptions,
};
use ironflow_core::error::OperationError;
use ironflow_core::operation::{Operation, OperationContext, TypedOperation};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use tokio_stream::StreamExt;
use crate::containers::DockerRef;
use crate::helpers::{docker_error, to_value};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ImageListEntry {
pub id: String,
pub repo_tags: Vec<String>,
pub size: i64,
pub created: i64,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ImageListOutput {
pub images: Vec<ImageListEntry>,
}
pub struct ImageList {
docker: Docker,
all: bool,
filters: HashMap<String, Vec<String>>,
}
impl ImageList {
pub fn new(client: impl Into<DockerRef>) -> Self {
Self {
docker: client.into().0,
all: false,
filters: HashMap::new(),
}
}
pub fn all(mut self) -> Self {
self.all = true;
self
}
pub fn filter(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.filters
.entry(key.into())
.or_default()
.push(value.into());
self
}
pub async fn run(&self, _ctx: &OperationContext) -> Result<ImageListOutput, OperationError> {
let options = ListImagesOptions {
all: self.all,
filters: Some(self.filters.clone()),
..Default::default()
};
let images = self
.docker
.list_images(Some(options))
.await
.map_err(docker_error)?;
let entries = images
.into_iter()
.map(|i| ImageListEntry {
id: i.id,
repo_tags: i.repo_tags,
size: i.size,
created: i.created,
})
.collect();
Ok(ImageListOutput { images: entries })
}
}
#[async_trait]
impl Operation for ImageList {
fn kind(&self) -> &str {
"docker"
}
async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
to_value(&self.run(ctx).await?)
}
fn input(&self) -> Option<Value> {
Some(serde_json::json!({
"operation": "image_list",
"all": self.all,
}))
}
}
impl TypedOperation for ImageList {
type Output = ImageListOutput;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ImagePullOutput {
pub image: String,
}
pub struct ImagePull {
docker: Docker,
image: String,
}
impl ImagePull {
pub fn new(client: impl Into<DockerRef>, image: impl Into<String>) -> Self {
Self {
docker: client.into().0,
image: image.into(),
}
}
pub async fn run(&self, _ctx: &OperationContext) -> Result<ImagePullOutput, OperationError> {
let options = CreateImageOptions {
from_image: Some(self.image.clone()),
..Default::default()
};
let mut stream = self.docker.create_image(Some(options), None, None);
while let Some(result) = stream.next().await {
result.map_err(docker_error)?;
}
Ok(ImagePullOutput {
image: self.image.clone(),
})
}
}
#[async_trait]
impl Operation for ImagePull {
fn kind(&self) -> &str {
"docker"
}
async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
to_value(&self.run(ctx).await?)
}
fn input(&self) -> Option<Value> {
Some(serde_json::json!({
"operation": "image_pull",
"image": self.image,
}))
}
}
impl TypedOperation for ImagePull {
type Output = ImagePullOutput;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ImagePushOutput {
pub image: String,
}
pub struct ImagePush {
docker: Docker,
image: String,
tag: Option<String>,
}
impl ImagePush {
pub fn new(client: impl Into<DockerRef>, image: impl Into<String>) -> Self {
Self {
docker: client.into().0,
image: image.into(),
tag: None,
}
}
pub fn tag(mut self, tag: impl Into<String>) -> Self {
self.tag = Some(tag.into());
self
}
pub async fn run(&self, _ctx: &OperationContext) -> Result<ImagePushOutput, OperationError> {
let options = PushImageOptions {
tag: Some(self.tag.clone().unwrap_or_else(|| "latest".to_string())),
platform: None,
};
let mut stream = self.docker.push_image(&self.image, Some(options), None);
while let Some(result) = stream.next().await {
result.map_err(docker_error)?;
}
Ok(ImagePushOutput {
image: self.image.clone(),
})
}
}
#[async_trait]
impl Operation for ImagePush {
fn kind(&self) -> &str {
"docker"
}
async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
to_value(&self.run(ctx).await?)
}
fn input(&self) -> Option<Value> {
Some(serde_json::json!({
"operation": "image_push",
"image": self.image,
}))
}
}
impl TypedOperation for ImagePush {
type Output = ImagePushOutput;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ImageInspectOutput {
pub data: Value,
}
pub struct ImageInspect {
docker: Docker,
image: String,
}
impl ImageInspect {
pub fn new(client: impl Into<DockerRef>, image: impl Into<String>) -> Self {
Self {
docker: client.into().0,
image: image.into(),
}
}
pub async fn run(&self, _ctx: &OperationContext) -> Result<ImageInspectOutput, OperationError> {
let response = self
.docker
.inspect_image(&self.image)
.await
.map_err(docker_error)?;
let data = serde_json::to_value(&response).map_err(|e| OperationError::External {
origin: "docker".to_string(),
message: e.to_string(),
})?;
Ok(ImageInspectOutput { data })
}
}
#[async_trait]
impl Operation for ImageInspect {
fn kind(&self) -> &str {
"docker"
}
async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
to_value(&self.run(ctx).await?)
}
fn input(&self) -> Option<Value> {
Some(serde_json::json!({
"operation": "image_inspect",
"image": self.image,
}))
}
}
impl TypedOperation for ImageInspect {
type Output = ImageInspectOutput;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ImageRemoveOutput {
pub image: String,
}
pub struct ImageRemove {
docker: Docker,
image: String,
force: bool,
no_prune: bool,
}
impl ImageRemove {
pub fn new(client: impl Into<DockerRef>, image: impl Into<String>) -> Self {
Self {
docker: client.into().0,
image: image.into(),
force: false,
no_prune: false,
}
}
pub fn force(mut self) -> Self {
self.force = true;
self
}
pub fn no_prune(mut self) -> Self {
self.no_prune = true;
self
}
pub async fn run(&self, _ctx: &OperationContext) -> Result<ImageRemoveOutput, OperationError> {
let options = RemoveImageOptions {
force: self.force,
noprune: self.no_prune,
platforms: None,
};
self.docker
.remove_image(&self.image, Some(options), None)
.await
.map_err(docker_error)?;
Ok(ImageRemoveOutput {
image: self.image.clone(),
})
}
}
#[async_trait]
impl Operation for ImageRemove {
fn kind(&self) -> &str {
"docker"
}
async fn execute(&self, ctx: &OperationContext) -> Result<Value, OperationError> {
to_value(&self.run(ctx).await?)
}
fn input(&self) -> Option<Value> {
Some(serde_json::json!({
"operation": "image_remove",
"image": self.image,
}))
}
}
impl TypedOperation for ImageRemove {
type Output = ImageRemoveOutput;
}