This commit is contained in:
@@ -0,0 +1,352 @@
|
||||
//! Client de l'API du daemon de containers et opérations de cycle de vie.
|
||||
|
||||
use std::{collections::HashMap, path::Path, time::Duration};
|
||||
|
||||
use bollard::{
|
||||
Docker, body_try_stream,
|
||||
container::LogOutput,
|
||||
errors::Error as BollardError,
|
||||
exec::{CreateExecOptions, StartExecOptions, StartExecResults},
|
||||
models::{BuildInfo, ContainerCreateBody, NetworkCreateRequest, NetworkDisconnectRequest},
|
||||
query_parameters::{
|
||||
BuildImageOptions, CreateContainerOptions, RemoveContainerOptions, RemoveImageOptions,
|
||||
StartContainerOptions, StopContainerOptions,
|
||||
},
|
||||
};
|
||||
use futures_util::StreamExt;
|
||||
|
||||
use crate::{
|
||||
consts::{DEFAULT_COMMAND_TIMEOUT, DEFAULT_ENDPOINT, EXIT_CODE_ATTEMPTS, EXIT_CODE_DELAY},
|
||||
context::tar_directory_stream,
|
||||
errors::ContainerError,
|
||||
exec::ExecOutput,
|
||||
};
|
||||
|
||||
/// Un client de l'API du daemon de containers.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ContainerRuntime {
|
||||
docker: Docker,
|
||||
endpoint: String,
|
||||
timeout: Duration,
|
||||
}
|
||||
|
||||
impl ContainerRuntime {
|
||||
/// Se connecte au daemon désigné par `DOCKER_HOST`, ou au socket local par défaut.
|
||||
pub fn connect() -> Result<Self, ContainerError> {
|
||||
let endpoint =
|
||||
std::env::var("DOCKER_HOST").unwrap_or_else(|_| String::from(DEFAULT_ENDPOINT));
|
||||
|
||||
Ok(Self {
|
||||
docker: Docker::connect_with_defaults()?,
|
||||
endpoint,
|
||||
timeout: DEFAULT_COMMAND_TIMEOUT,
|
||||
})
|
||||
}
|
||||
|
||||
/// Remplace le timeout appliqué aux opérations de build/run/stop/remove.
|
||||
pub fn with_timeout(mut self, timeout: Duration) -> Self {
|
||||
self.timeout = timeout;
|
||||
self
|
||||
}
|
||||
|
||||
/// Endpoint du daemon, pour les logs.
|
||||
pub fn endpoint(&self) -> &str {
|
||||
&self.endpoint
|
||||
}
|
||||
|
||||
pub fn timeout(&self) -> Duration {
|
||||
self.timeout
|
||||
}
|
||||
|
||||
/// Vérifie que le daemon est joignable et répond.
|
||||
pub async fn available(&self) -> bool {
|
||||
self.docker.ping().await.is_ok()
|
||||
}
|
||||
|
||||
/// Exécute une requête du daemon en appliquant le timeout des commandes.
|
||||
async fn request<T>(
|
||||
&self,
|
||||
operation: &str,
|
||||
request: impl Future<Output = Result<T, BollardError>>,
|
||||
) -> Result<T, ContainerError> {
|
||||
match tokio::time::timeout(self.timeout, request).await {
|
||||
Ok(Ok(value)) => Ok(value),
|
||||
Ok(Err(source)) => Err(ContainerError::Request { source }),
|
||||
Err(_) => Err(ContainerError::Timeout {
|
||||
operation: String::from(operation),
|
||||
timeout: self.timeout,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Construit l'image `tag` depuis le contexte `context_dir`.
|
||||
pub(crate) async fn build_image(
|
||||
&self,
|
||||
context_dir: &Path,
|
||||
dockerfile: &str,
|
||||
tag: &str,
|
||||
build_args: &HashMap<String, String>,
|
||||
) -> Result<(), ContainerError> {
|
||||
let options = BuildImageOptions {
|
||||
dockerfile: String::from(dockerfile),
|
||||
t: Some(String::from(tag)),
|
||||
buildargs: Some(build_args.clone()),
|
||||
rm: true,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let context = tar_directory_stream(context_dir);
|
||||
let mut stream = self
|
||||
.docker
|
||||
.build_image(options, None, Some(body_try_stream(context)));
|
||||
|
||||
// Le flux porte la progression du build : la dernière erreur signalée par le
|
||||
// daemon fait échouer l'opération.
|
||||
let consume = async {
|
||||
let mut failure = None;
|
||||
|
||||
while let Some(info) = stream.next().await {
|
||||
let info: BuildInfo = info?;
|
||||
|
||||
if let Some(message) = info.error_detail.and_then(|detail| detail.message) {
|
||||
failure = Some(message);
|
||||
}
|
||||
}
|
||||
|
||||
Ok::<_, BollardError>(failure)
|
||||
};
|
||||
|
||||
let failure = match tokio::time::timeout(self.timeout, consume).await {
|
||||
Ok(Ok(failure)) => failure,
|
||||
Ok(Err(source)) => return Err(ContainerError::Request { source }),
|
||||
Err(_) => {
|
||||
return Err(ContainerError::Timeout {
|
||||
operation: format!("build image `{tag}`"),
|
||||
timeout: self.timeout,
|
||||
});
|
||||
}
|
||||
};
|
||||
|
||||
match failure {
|
||||
Some(message) => Err(ContainerError::Build { message }),
|
||||
None => Ok(()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Crée un container, sans le démarrer.
|
||||
pub(crate) async fn create_container(
|
||||
&self,
|
||||
name: &str,
|
||||
body: ContainerCreateBody,
|
||||
) -> Result<(), ContainerError> {
|
||||
let options = CreateContainerOptions {
|
||||
name: Some(String::from(name)),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.request(
|
||||
"create container",
|
||||
self.docker.create_container(Some(options), body),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn start_container(&self, name: &str) -> Result<(), ContainerError> {
|
||||
self.request(
|
||||
"start container",
|
||||
self.docker
|
||||
.start_container(name, None::<StartContainerOptions>),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Exécute une commande dans un container et renvoie sa sortie.
|
||||
pub(crate) async fn exec(
|
||||
&self,
|
||||
container: &str,
|
||||
cmd: &[&str],
|
||||
user: Option<&str>,
|
||||
timeout: Duration,
|
||||
) -> Result<ExecOutput, ContainerError> {
|
||||
let config = CreateExecOptions {
|
||||
cmd: Some(cmd.iter().map(|arg| String::from(*arg)).collect()),
|
||||
user: user.map(String::from),
|
||||
attach_stdout: Some(true),
|
||||
attach_stderr: Some(true),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let created = self
|
||||
.request("create exec", self.docker.create_exec(container, config))
|
||||
.await?;
|
||||
|
||||
let started = self
|
||||
.request(
|
||||
"start exec",
|
||||
self.docker
|
||||
.start_exec(&created.id, None::<StartExecOptions>),
|
||||
)
|
||||
.await?;
|
||||
|
||||
let StartExecResults::Attached { mut output, .. } = started else {
|
||||
return Err(ContainerError::Unexpected(String::from(
|
||||
"the daemon detached an exec that was not requested as detached",
|
||||
)));
|
||||
};
|
||||
|
||||
let mut stdout = String::new();
|
||||
let mut stderr = String::new();
|
||||
|
||||
// Le timeout couvre l'exécution de la commande elle-même, pas seulement sa
|
||||
// mise en place.
|
||||
let collect = async {
|
||||
while let Some(message) = output.next().await {
|
||||
match message? {
|
||||
LogOutput::StdOut { message } | LogOutput::Console { message } => {
|
||||
stdout.push_str(&String::from_utf8_lossy(&message));
|
||||
}
|
||||
LogOutput::StdErr { message } => {
|
||||
stderr.push_str(&String::from_utf8_lossy(&message));
|
||||
}
|
||||
LogOutput::StdIn { .. } => {}
|
||||
}
|
||||
}
|
||||
|
||||
Ok::<_, BollardError>(())
|
||||
};
|
||||
|
||||
match tokio::time::timeout(timeout, collect).await {
|
||||
Ok(Ok(())) => {}
|
||||
Ok(Err(source)) => return Err(ContainerError::Request { source }),
|
||||
Err(_) => {
|
||||
// Le processus continue de tourner dans le container : sans arrêt, il
|
||||
// consommerait CPU et mémoire et pourrait encore modifier le workspace
|
||||
// pendant toute la durée de vie de la sandbox. L'API n'offre pas de
|
||||
// moyen de tuer un exec, on arrête donc le container qui le porte.
|
||||
let _ = self.stop_container(container).await;
|
||||
|
||||
return Err(ContainerError::Timeout {
|
||||
operation: format!("exec {}", cmd.join(" ")),
|
||||
timeout,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
// Le code de sortie n'est pas forcément publié au moment où le flux se ferme :
|
||||
// on laisse au daemon le temps de le renseigner, sans bloquer indéfiniment.
|
||||
// Sans cela, une commande réussie peut être rapportée en échec (code -1) et
|
||||
// l'outil renvoie une erreur au modèle à la place du contenu.
|
||||
let mut inspected = self
|
||||
.request("inspect exec", self.docker.inspect_exec(&created.id))
|
||||
.await?;
|
||||
|
||||
for _ in 0..EXIT_CODE_ATTEMPTS {
|
||||
if !inspected.running.unwrap_or(false) {
|
||||
break;
|
||||
}
|
||||
|
||||
tokio::time::sleep(EXIT_CODE_DELAY).await;
|
||||
|
||||
inspected = self
|
||||
.request("inspect exec", self.docker.inspect_exec(&created.id))
|
||||
.await?;
|
||||
}
|
||||
|
||||
Ok(ExecOutput {
|
||||
status: inspected.exit_code.unwrap_or(-1) as i32,
|
||||
stdout,
|
||||
stderr,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) async fn stop_container(&self, name: &str) -> Result<(), ContainerError> {
|
||||
let options = StopContainerOptions {
|
||||
t: Some(1),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.request(
|
||||
"stop container",
|
||||
self.docker.stop_container(name, Some(options)),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Supprime un container, ses volumes anonymes et le processus qui y tourne.
|
||||
pub(crate) async fn remove_container(&self, name: &str) -> Result<(), ContainerError> {
|
||||
let options = RemoveContainerOptions {
|
||||
force: true,
|
||||
v: true,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.request(
|
||||
"remove container",
|
||||
self.docker.remove_container(name, Some(options)),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn remove_image(&self, tag: &str) -> Result<(), ContainerError> {
|
||||
let options = RemoveImageOptions {
|
||||
force: true,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.request(
|
||||
"remove image",
|
||||
self.docker.remove_image(tag, Some(options), None),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Crée le réseau dédié d'une sandbox.
|
||||
pub(crate) async fn create_network(&self, name: &str) -> Result<(), ContainerError> {
|
||||
let request = NetworkCreateRequest {
|
||||
name: String::from(name),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.request("create network", self.docker.create_network(request))
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Détache un container d'un réseau, ce qui coupe sa connectivité.
|
||||
pub(crate) async fn disconnect_network(
|
||||
&self,
|
||||
network: &str,
|
||||
container: &str,
|
||||
) -> Result<(), ContainerError> {
|
||||
let request = NetworkDisconnectRequest {
|
||||
container: String::from(container),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
self.request(
|
||||
"disconnect network",
|
||||
self.docker.disconnect_network(network, request),
|
||||
)
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub(crate) async fn remove_network(&self, name: &str) -> Result<(), ContainerError> {
|
||||
self.request("remove network", self.docker.remove_network(name))
|
||||
.await?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user