Files
2026-09-20 13:10:05 +00:00

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(())
}
}