377 lines
12 KiB
Rust
377 lines
12 KiB
Rust
//! 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, UploadToContainerOptions,
|
|
},
|
|
};
|
|
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 upload_directory(
|
|
&self,
|
|
container: &str,
|
|
destination: &str,
|
|
directory: &Path,
|
|
) -> Result<(), ContainerError> {
|
|
let options = UploadToContainerOptions {
|
|
path: String::from(destination),
|
|
..Default::default()
|
|
};
|
|
|
|
self.request(
|
|
"upload workspace",
|
|
self.docker.upload_to_container(
|
|
container,
|
|
Some(options),
|
|
body_try_stream(tar_directory_stream(directory)),
|
|
),
|
|
)
|
|
.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(())
|
|
}
|
|
}
|