use async_trait::async_trait;
use boatramp_core::compute::{
Artifact, BackendError, Capabilities, ComputeBackend, ComputeSpec, Endpoint, Health, Instance,
InstanceHandle, IsolationClass, LaunchRequest, RestartPolicy, Scheme,
};
use bollard::container::{
Config, CreateContainerOptions, RemoveContainerOptions, StopContainerOptions,
};
use bollard::image::CreateImageOptions;
use bollard::models::{HostConfig, RestartPolicy as DockerRestartPolicy, RestartPolicyNameEnum};
use bollard::Docker;
use futures::StreamExt;
pub struct DockerBackend {
docker: Docker,
}
impl DockerBackend {
pub fn connect() -> Result<Self, BackendError> {
let docker = Docker::connect_with_defaults()
.map_err(|e| BackendError::Other(format!("connect to docker: {e}")))?;
Ok(Self { docker })
}
pub fn with_client(docker: Docker) -> Self {
Self { docker }
}
pub async fn reachable(&self) -> bool {
self.docker.ping().await.is_ok()
}
}
fn container_name(workload: &str, replica: u32) -> String {
format!("boatramp-{workload}-{replica}")
}
fn encode_ref(name: &str, ip: &str, port: u16) -> String {
format!("{name}@{ip}:{port}")
}
fn decode_ref(s: &str) -> Option<(String, String, u16)> {
let (name, rest) = s.split_once('@')?;
let (ip, port) = rest.rsplit_once(':')?;
Some((name.to_string(), ip.to_string(), port.parse().ok()?))
}
fn restart_policy(policy: RestartPolicy) -> DockerRestartPolicy {
let name = match policy {
RestartPolicy::Never => RestartPolicyNameEnum::NO,
RestartPolicy::OnFailure => RestartPolicyNameEnum::ON_FAILURE,
RestartPolicy::Always => RestartPolicyNameEnum::ALWAYS,
};
DockerRestartPolicy {
name: Some(name),
maximum_retry_count: None,
}
}
const MAX_PIDS: i64 = 512;
fn hardened_host_config(mem_mib: u32, vcpus: u32, restart: RestartPolicy) -> HostConfig {
let tmpfs = std::collections::HashMap::from([
("/tmp".to_string(), "rw,noexec,nosuid,size=64m".to_string()),
("/run".to_string(), "rw,noexec,nosuid,size=16m".to_string()),
]);
HostConfig {
memory: Some(i64::from(mem_mib) * 1024 * 1024),
nano_cpus: Some(i64::from(vcpus.max(1)) * 1_000_000_000),
restart_policy: Some(restart_policy(restart)),
security_opt: Some(vec!["no-new-privileges:true".to_string()]),
cap_drop: Some(vec!["ALL".to_string()]),
readonly_rootfs: Some(true),
tmpfs: Some(tmpfs),
pids_limit: Some(MAX_PIDS),
..Default::default()
}
}
#[async_trait]
impl ComputeBackend for DockerBackend {
fn id(&self) -> &'static str {
"docker"
}
fn capabilities(&self) -> Capabilities {
Capabilities {
isolation: IsolationClass::Container,
scale_to_zero: false,
persistent_volumes: false,
max_vcpus: None,
max_mem_mib: None,
}
}
async fn materialize(&self, spec: &ComputeSpec) -> Result<Artifact, BackendError> {
let reference = spec.rootfs.clone();
let options = CreateImageOptions {
from_image: reference.clone(),
..Default::default()
};
let mut pull = self.docker.create_image(Some(options), None, None);
while let Some(step) = pull.next().await {
step.map_err(|e| BackendError::Materialize(format!("pull {reference}: {e}")))?;
}
Ok(Artifact::Image { reference })
}
async fn launch(&self, req: &LaunchRequest) -> Result<Instance, BackendError> {
let reference = match &req.artifact {
Artifact::Image { reference } => reference.clone(),
_ => {
return Err(BackendError::Launch(
"docker backend requires an Image artifact".into(),
))
}
};
let name = container_name(&req.workload, req.replica);
let env: Vec<String> = req
.spec
.env
.iter()
.map(|(k, v)| format!("{k}={v}"))
.collect();
let host_config = hardened_host_config(req.spec.mem_mib, req.spec.vcpus, req.spec.restart);
let config = Config {
image: Some(reference),
cmd: Some(req.spec.entrypoint.clone()),
env: Some(env),
host_config: Some(host_config),
..Default::default()
};
let _ = self
.docker
.remove_container(
&name,
Some(RemoveContainerOptions {
force: true,
..Default::default()
}),
)
.await;
let created = self
.docker
.create_container(
Some(CreateContainerOptions {
name: name.clone(),
platform: None,
}),
config,
)
.await
.map_err(|e| BackendError::Launch(format!("create {name}: {e}")))?;
self.docker
.start_container::<String>(&created.id, None)
.await
.map_err(|e| BackendError::Launch(format!("start {name}: {e}")))?;
let ip = self.container_ip(&created.id).await?;
let port = req.spec.port;
Ok(Instance {
handle: InstanceHandle {
workload: req.workload.clone(),
replica: req.replica,
backend_ref: encode_ref(&name, &ip, port),
},
endpoint: Endpoint {
scheme: Scheme::Http,
host: ip,
port,
},
})
}
async fn stop(&self, handle: &InstanceHandle) -> Result<(), BackendError> {
let name = decode_ref(&handle.backend_ref)
.map(|(n, _, _)| n)
.unwrap_or_else(|| container_name(&handle.workload, handle.replica));
let _ = self
.docker
.stop_container(&name, None::<StopContainerOptions>)
.await;
self.docker
.remove_container(
&name,
Some(RemoveContainerOptions {
force: true,
..Default::default()
}),
)
.await
.map_err(|e| BackendError::Stop(format!("remove {name}: {e}")))?;
Ok(())
}
async fn health(&self, handle: &InstanceHandle) -> Result<Health, BackendError> {
let name = match decode_ref(&handle.backend_ref) {
Some((n, _, _)) => n,
None => container_name(&handle.workload, handle.replica),
};
let info = match self.docker.inspect_container(&name, None).await {
Ok(info) => info,
Err(_) => return Ok(Health::Unhealthy),
};
let running = info.state.and_then(|s| s.running).unwrap_or(false);
Ok(if running {
Health::Healthy
} else {
Health::Unhealthy
})
}
}
impl DockerBackend {
async fn container_ip(&self, id: &str) -> Result<String, BackendError> {
let info = self
.docker
.inspect_container(id, None)
.await
.map_err(|e| BackendError::Launch(format!("inspect {id}: {e}")))?;
let networks = info
.network_settings
.ok_or_else(|| BackendError::Launch("container has no network settings".into()))?;
if let Some(ip) = networks.ip_address.filter(|s| !s.is_empty()) {
return Ok(ip);
}
if let Some(nets) = networks.networks {
for net in nets.values() {
if let Some(ip) = net.ip_address.as_ref().filter(|s| !s.is_empty()) {
return Ok(ip.clone());
}
}
}
Err(BackendError::Launch("container has no IP address".into()))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn name_and_ref_round_trip() {
assert_eq!(container_name("web", 0), "boatramp-web-0");
let r = encode_ref("boatramp-web-0", "172.17.0.3", 8080);
assert_eq!(r, "boatramp-web-0@172.17.0.3:8080");
assert_eq!(
decode_ref(&r),
Some(("boatramp-web-0".to_string(), "172.17.0.3".to_string(), 8080))
);
assert_eq!(decode_ref("garbage"), None);
}
#[test]
fn host_config_is_hardened_by_default() {
let hc = hardened_host_config(256, 2, RestartPolicy::Never);
assert_eq!(hc.memory, Some(256 * 1024 * 1024));
assert_eq!(hc.nano_cpus, Some(2_000_000_000));
assert_eq!(hc.pids_limit, Some(MAX_PIDS));
assert_eq!(
hc.security_opt.as_deref(),
Some(["no-new-privileges:true".to_string()].as_slice())
);
assert_eq!(hc.cap_drop.as_deref(), Some(["ALL".to_string()].as_slice()));
assert_eq!(hc.readonly_rootfs, Some(true));
let tmpfs = hc.tmpfs.expect("tmpfs mounts for a read-only rootfs");
assert!(tmpfs.get("/tmp").is_some_and(|o| o.contains("noexec")));
assert!(tmpfs.contains_key("/run"));
assert_eq!(
hardened_host_config(64, 0, RestartPolicy::Never).nano_cpus,
Some(1_000_000_000)
);
}
#[test]
fn restart_policy_maps_to_docker() {
assert_eq!(
restart_policy(RestartPolicy::Always).name,
Some(RestartPolicyNameEnum::ALWAYS)
);
assert_eq!(
restart_policy(RestartPolicy::OnFailure).name,
Some(RestartPolicyNameEnum::ON_FAILURE)
);
assert_eq!(
restart_policy(RestartPolicy::Never).name,
Some(RestartPolicyNameEnum::NO)
);
}
}