use std::collections::HashMap;
#[cfg(unix)]
use std::os::unix::process::ExitStatusExt as _;
#[cfg(windows)]
use std::os::windows::process::ExitStatusExt as _;
use std::path::PathBuf;
use std::process::ExitStatus;
use std::time::Duration;
use bollard::Docker;
use bollard::models::ContainerWaitResponse;
use bollard::models::TaskState;
use bollard::query_parameters::InspectContainerOptions;
use bollard::query_parameters::ListTasksOptions;
use bollard::query_parameters::LogsOptionsBuilder;
use bollard::query_parameters::WaitContainerOptions;
mod builder;
pub use builder::Builder;
use crankshaft_events::Event;
use futures::StreamExt;
use tokio::fs::File;
use tokio::time::sleep;
use tracing::debug;
use tracing::info;
use tracing::trace;
use crate::Error;
use crate::EventOptions;
use crate::Result;
use crate::container::ExecutionResult;
use crate::container::write_logs;
pub struct Service {
client: Docker,
id: String,
stdout: Option<PathBuf>,
stderr: Option<PathBuf>,
}
impl Service {
pub fn new(
client: Docker,
id: String,
stdout: Option<PathBuf>,
stderr: Option<PathBuf>,
) -> Self {
Self {
client,
id,
stdout,
stderr,
}
}
pub fn id(&self) -> &str {
&self.id
}
pub async fn run(
&self,
task_name: &str,
events: Option<EventOptions>,
) -> Result<ExecutionResult> {
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,
};
let (image, container_id, exit_code) = loop {
trace!(
"polling tasks for service `{id}` (task `{task_name}`)",
id = self.id
);
let tasks = self
.client
.list_tasks(Some(ListTasksOptions {
filters: Some(HashMap::from_iter([(
String::from("service"),
vec![self.id.to_owned()],
)])),
}))
.await
.map_err(Error::Docker)?;
if tasks.is_empty() {
sleep(Duration::from_millis(100)).await;
continue;
}
assert_eq!(
tasks.len(),
1,
"Docker service task count should always be 1"
);
let task = tasks.into_iter().next().unwrap();
let image = task
.spec
.as_ref()
.and_then(|spec| spec.container_spec.as_ref())
.and_then(|spec| spec.image.clone())
.expect("Docker reported a container without a state");
let status = task.status.ok_or_else(|| {
Error::Message("Docker daemon reported a task with no status".into())
})?;
match status.state {
Some(TaskState::NEW)
| Some(TaskState::PENDING)
| Some(TaskState::ALLOCATED)
| Some(TaskState::ASSIGNED)
| Some(TaskState::ACCEPTED)
| Some(TaskState::READY)
| Some(TaskState::PREPARING)
| Some(TaskState::STARTING)
| None => {
trace!(
"task has not yet started for service `{id}` (task `{task_name}`)",
id = self.id
);
sleep(Duration::from_secs(1)).await;
}
Some(TaskState::RUNNING) | Some(TaskState::COMPLETE) | Some(TaskState::FAILED) => {
let container_status = status.container_status.ok_or_else(|| {
Error::Message(
"Docker daemon reported a task with no container status".into(),
)
})?;
let container_id = container_status.container_id.ok_or_else(|| {
Error::Message("Docker reported a task with no container id".into())
})?;
if let Some(events) = &events {
events
.sender
.send(Event::TaskContainerCreated {
id: events.task_id,
container: container_id.clone(),
})
.ok();
}
info!(
"service `{id}` (task `{task_name}`) has started container `{container_id}",
id = self.id
);
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(
&container_id,
Some(
LogsOptionsBuilder::new()
.stdout(stdout_enabled)
.stderr(stderr_enabled)
.follow(true)
.build(),
),
);
write_logs(logs, stdout, stderr, events.as_ref()).await?;
}
if status.state == Some(TaskState::RUNNING) {
let mut wait_stream = self
.client
.wait_container(&container_id, None::<WaitContainerOptions>);
match wait_stream.next().await {
Some(Ok(ContainerWaitResponse {
status_code: code, ..
}))
| Some(Err(bollard::errors::Error::DockerContainerWaitError {
code,
..
})) => {
break (image, container_id, code);
}
Some(Err(e)) => return Err(e.into()),
None => {
let container = self
.client
.inspect_container(
&container_id,
None::<InspectContainerOptions>,
)
.await
.map_err(Error::Docker)?;
break (
image,
container_id,
container
.state
.ok_or_else(|| {
Error::Message(
"Docker reported a container without a state"
.into(),
)
})?
.exit_code
.ok_or_else(|| {
Error::Message(
"Docker reported a finished contained without an \
exit code"
.into(),
)
})?,
);
}
}
} else {
break (
image,
container_id,
container_status.exit_code.ok_or_else(|| {
Error::Message(format!(
"Docker reported a {kind} task with no exit code",
kind = if status.state == Some(TaskState::FAILED) {
"failed"
} else {
"completed"
}
))
})?,
);
}
}
Some(TaskState::SHUTDOWN)
| Some(TaskState::REJECTED)
| Some(TaskState::ORPHANED)
| Some(TaskState::REMOVE) => {
return Err(Error::Message(format!(
"Docker task failed: {msg}",
msg = status
.err
.as_deref()
.or(status.message.as_deref())
.unwrap_or("no error message was provided by the Docker daemon")
)));
}
}
};
#[cfg(unix)]
let status = ExitStatus::from_raw((exit_code as i32) << 8);
#[cfg(windows)]
let status = ExitStatus::from_raw(exit_code as u32);
info!(
"container `{container_id}` for service `{id}` (task `{task_name}`) has exited with \
{status}",
id = self.id
);
if let Some(events) = &events {
events
.sender
.send(Event::TaskContainerExited {
id: events.task_id,
container: container_id,
exit_status: status,
})
.ok();
}
Ok(ExecutionResult { image, status })
}
pub async fn delete(&self) -> Result<()> {
debug!("deleting Docker service `{id}`", id = self.id);
self.client
.delete_service(&self.id)
.await
.map_err(Error::Docker)?;
Ok(())
}
}