//! 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 { 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( &self, operation: &str, request: impl Future>, ) -> Result { 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, ) -> 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::), ) .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 { 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::), ) .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(()) } }