use std::future::Future;
#[cfg(unix)]
use std::os::unix::process::ExitStatusExt as _;
#[cfg(windows)]
use std::os::windows::process::ExitStatusExt as _;
use std::path::Path;
use std::path::PathBuf;
use std::process::ExitStatus;
use std::time::Duration;
use bollard::Docker;
use bollard::container::LogOutput;
use bollard::models::ContainerWaitResponse;
use bollard::query_parameters::InspectContainerOptions;
use bollard::query_parameters::LogsOptionsBuilder;
use bollard::query_parameters::RemoveContainerOptions;
use bollard::query_parameters::StartContainerOptions;
use bollard::query_parameters::WaitContainerOptions;
use crankshaft_events::Event;
use futures::Stream;
use tokio::fs::File;
use tokio::io::AsyncWriteExt;
use tokio::pin;
use tokio_retry2::Retry;
use tokio_retry2::RetryError;
use tokio_retry2::strategy::ExponentialBackoff;
use tokio_retry2::strategy::jitter;
use tokio_stream::StreamExt as _;
use tracing::debug;
use tracing::info;
use crate::Error;
use crate::EventOptions;
use crate::Result;
mod builder;
pub use builder::Builder;
fn default_retry_strategy() -> impl Iterator<Item = Duration> {
ExponentialBackoff::from_millis(50)
.factor(2)
.max_delay_millis(1000)
.map(jitter)
.take(5)
}
fn is_retryable(err: bollard::errors::Error) -> RetryError<bollard::errors::Error> {
match err {
bollard::errors::Error::IOError { .. } |
bollard::errors::Error::HyperResponseError { .. } => RetryError::transient(err),
bollard::errors::Error::DockerResponseServerError { status_code: 409, .. } => RetryError::transient(err),
_ => {
RetryError::permanent(err)
}
}
}
async fn default_retry<F, Fut, T>(op: F) -> std::result::Result<T, bollard::errors::Error>
where
F: Fn() -> Fut,
Fut: Future<Output = std::result::Result<T, bollard::errors::Error>>,
{
let action = || async { op().await.map_err(is_retryable) };
Retry::spawn(default_retry_strategy(), action).await
}
pub(crate) async fn write_logs(
logs: impl Stream<Item = std::result::Result<LogOutput, bollard::errors::Error>>,
mut stdout: Option<(&Path, File)>,
mut stderr: Option<(&Path, File)>,
events: Option<&EventOptions>,
) -> Result<()> {
pin!(logs);
while let Some(result) = logs.next().await {
let output = result.map_err(Error::Docker)?;
match output {
LogOutput::StdOut { message } => {
if let Some((path, stdout)) = &mut stdout {
stdout.write(&message).await.map_err(|e| {
Error::Message(format!(
"failed to write to stdout file `{path}`: {e}",
path = path.display()
))
})?;
}
if let Some(events) = events
&& events.user_config.send_stdout
{
events
.sender
.send(Event::TaskStdout {
id: events.task_id,
message,
})
.ok();
}
}
LogOutput::StdErr { message } => {
if let Some((path, stderr)) = &mut stderr {
stderr.write(&message).await.map_err(|e| {
Error::Message(format!(
"failed to write to stderr file `{path}`: {e}",
path = path.display()
))
})?;
}
if let Some(events) = &events
&& events.user_config.send_stderr
{
events
.sender
.send(Event::TaskStderr {
id: events.task_id,
message,
})
.ok();
}
}
_ => {}
}
}
Ok(())
}
pub struct ExecutionResult {
pub image: String,
pub status: ExitStatus,
}
pub struct Container {
client: Docker,
name: String,
stdout: Option<PathBuf>,
stderr: Option<PathBuf>,
}
impl Container {
pub fn new(
client: Docker,
name: String,
stdout: Option<PathBuf>,
stderr: Option<PathBuf>,
) -> Self {
Self {
client,
name,
stdout,
stderr,
}
}
pub fn name(&self) -> &str {
&self.name
}
pub async fn run(
&self,
task_name: &str,
events: Option<EventOptions>,
) -> Result<ExecutionResult> {
if let Some(events) = &events {
events
.sender
.send(Event::TaskContainerCreated {
id: events.task_id,
container: self.name.clone(),
})
.ok();
}
info!(
"starting container `{name}` (task `{task_name}`)",
name = self.name
);
let inspect_response = default_retry(|| {
self.client
.inspect_container(&self.name, None::<InspectContainerOptions>)
})
.await
.map_err(Error::Docker)?;
default_retry(|| {
self.client
.start_container(&self.name, None::<StartContainerOptions>)
})
.await
.map_err(Error::Docker)?;
info!(
"container `{name}` (task `{task_name}`) has started",
name = self.name
);
if let Some(events) = &events
&& events.send_start
{
events
.sender
.send(Event::TaskStarted { id: events.task_id })
.ok();
}
let stdout_enabled =
self.stdout.is_some() || events.as_ref().is_some_and(|e| e.user_config.send_stdout);
let stderr_enabled =
self.stderr.is_some() || events.as_ref().is_some_and(|e| e.user_config.send_stderr);
if stdout_enabled || stderr_enabled {
let logs = self.client.logs(
&self.name,
Some(
LogsOptionsBuilder::new()
.stdout(stdout_enabled)
.stderr(stderr_enabled)
.follow(true)
.build(),
),
);
let stdout = match &self.stdout {
Some(path) => Some((
path.as_path(),
File::create(path).await.map_err(|e| {
Error::Message(format!(
"failed to create stdout file `{path}`: {e}",
path = path.display()
))
})?,
)),
None => None,
};
let stderr = match &self.stderr {
Some(path) => Some((
path.as_path(),
File::create(path).await.map_err(|e| {
Error::Message(format!(
"failed to create stderr file `{path}`: {e}",
path = path.display()
))
})?,
)),
None => None,
};
write_logs(logs, stdout, stderr, events.as_ref()).await?;
}
debug!(
"waiting for container `{name}` (task `{task_name}`) to exit",
name = self.name
);
let mut wait_stream = self
.client
.wait_container(&self.name, None::<WaitContainerOptions>);
let mut exit_code = None;
if let Some(result) = wait_stream.next().await {
match result {
Ok(ContainerWaitResponse {
status_code: code, ..
})
| Err(bollard::errors::Error::DockerContainerWaitError { code, .. }) => {
exit_code = Some(code);
}
Err(e) => return Err(e.into()),
}
}
if exit_code.is_none() {
let container = self
.client
.inspect_container(&self.name, None::<InspectContainerOptions>)
.await
.map_err(Error::Docker)?;
exit_code = Some(
container
.state
.expect("Docker reported a container without a state")
.exit_code
.expect("Docker reported a finished contained without an exit code"),
);
}
#[cfg(unix)]
let status = ExitStatus::from_raw((exit_code.unwrap() as i32) << 8);
#[cfg(windows)]
let status = ExitStatus::from_raw(exit_code.unwrap() as u32);
info!(
"container `{name}` (task `{task_name}`) has exited with {status}",
name = self.name
);
if let Some(events) = &events {
events
.sender
.send(Event::TaskContainerExited {
id: events.task_id,
container: self.name.clone(),
exit_status: status,
})
.ok();
}
Ok(ExecutionResult {
image: inspect_response
.config
.expect("Docker reported a container without a configuration")
.image
.expect("Docker reported a container without an image"),
status,
})
}
async fn remove_inner(&self, force: bool) -> Result<()> {
default_retry(|| {
self.client.remove_container(
&self.name,
Some(RemoveContainerOptions {
force,
..Default::default()
}),
)
})
.await
.map_err(Error::Docker)
}
pub async fn remove(&self) -> Result<()> {
debug!("removing container `{name}`", name = self.name);
self.remove_inner(false).await
}
pub async fn force_remove(&self) -> Result<()> {
debug!("force removing container `{name}`", name = self.name);
self.remove_inner(true).await
}
}