use async_trait::async_trait;
use boatramp_core::compute::{
Artifact, BackendError, Capabilities, ComputeBackend, ComputeSpec, Endpoint, Health, Instance,
InstanceHandle, IsolationClass, LaunchRequest, RestartPolicy, RootSource, Scheme, VolumeRef,
};
use bollard::container::{
Config, CreateContainerOptions, RemoveContainerOptions, StopContainerOptions,
};
use bollard::image::CreateImageOptions;
use bollard::models::{
HostConfig, Mount, MountTypeEnum, PortBinding, RestartPolicy as DockerRestartPolicy,
RestartPolicyNameEnum,
};
use bollard::Docker;
use futures::StreamExt;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum DockerEndpoint {
#[default]
Published,
Bridge,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum DockerVolumeMode {
#[default]
Named,
Bind,
}
fn docker_volume_name(name: &str) -> String {
format!("boatramp-{name}")
}
fn volume_dir(data_dir: &Path, name: &str) -> PathBuf {
data_dir.join("compute").join("volumes").join(name)
}
fn validate_volume(name: &str, mount: &str) -> Result<(), BackendError> {
use std::path::Component;
let name_ok = matches!(
Path::new(name).components().collect::<Vec<_>>().as_slice(),
[Component::Normal(_)]
);
if !name_ok {
return Err(BackendError::Launch(format!(
"invalid volume name {name:?}: must be a single path component"
)));
}
let m = Path::new(mount);
let mount_ok = m.is_absolute()
&& m.components()
.all(|c| matches!(c, Component::RootDir | Component::Normal(_)));
if !mount_ok {
return Err(BackendError::Launch(format!(
"invalid volume mount {mount:?}: must be an absolute path with no `..`"
)));
}
Ok(())
}
fn volume_mount(vol: &VolumeRef, mode: DockerVolumeMode, data_dir: &Path) -> Mount {
let (typ, source) = match mode {
DockerVolumeMode::Named => (MountTypeEnum::VOLUME, docker_volume_name(&vol.name)),
DockerVolumeMode::Bind => (
MountTypeEnum::BIND,
volume_dir(data_dir, &vol.name).display().to_string(),
),
};
Mount {
target: Some(vol.mount.clone()),
source: Some(source),
typ: Some(typ),
read_only: Some(false),
..Default::default()
}
}
pub struct DockerBackend {
docker: Docker,
endpoint: DockerEndpoint,
volume_mode: DockerVolumeMode,
data_dir: PathBuf,
writable_root_allowed: bool,
}
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,
endpoint: DockerEndpoint::default(),
volume_mode: DockerVolumeMode::default(),
data_dir: PathBuf::from("."),
writable_root_allowed: false,
})
}
pub fn with_client(docker: Docker) -> Self {
Self {
docker,
endpoint: DockerEndpoint::default(),
volume_mode: DockerVolumeMode::default(),
data_dir: PathBuf::from("."),
writable_root_allowed: false,
}
}
pub fn with_endpoint(mut self, endpoint: DockerEndpoint) -> Self {
self.endpoint = endpoint;
self
}
pub fn with_writable_root_allowed(mut self, allowed: bool) -> Self {
self.writable_root_allowed = allowed;
self
}
pub fn with_volume_mode(mut self, mode: DockerVolumeMode) -> Self {
self.volume_mode = mode;
self
}
pub fn with_data_dir(mut self, data_dir: impl Into<PathBuf>) -> Self {
self.data_dir = data_dir.into();
self
}
pub async fn reachable(&self) -> bool {
self.docker.ping().await.is_ok()
}
async fn stage_volumes(&self, spec: &ComputeSpec) -> Result<Vec<Mount>, BackendError> {
let mut mounts = Vec::with_capacity(spec.volumes.len());
for vol in &spec.volumes {
validate_volume(&vol.name, &vol.mount)?;
if self.volume_mode == DockerVolumeMode::Bind {
let dir = volume_dir(&self.data_dir, &vol.name);
tokio::fs::create_dir_all(&dir).await.map_err(|e| {
BackendError::Launch(format!("create volume {} dir: {e}", vol.name))
})?;
}
mounts.push(volume_mount(vol, self.volume_mode, &self.data_dir));
}
Ok(mounts)
}
}
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,
writable_root: bool,
) -> 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(!writable_root),
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: true,
max_vcpus: None,
max_mem_mib: None,
}
}
async fn materialize(&self, spec: &ComputeSpec) -> Result<Artifact, BackendError> {
let reference = match &spec.root {
RootSource::Image(reference) => reference.clone(),
RootSource::Tar(_) | RootSource::Rootfs(_) => {
return Err(BackendError::Materialize(
"docker backend requires an image reference (RootSource::Image)".into(),
))
}
};
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 port = req.spec.port;
let port_key = format!("{port}/tcp");
let writable_root = req.spec.writable_root && self.writable_root_allowed;
let mut host_config = hardened_host_config(
req.spec.mem_mib,
req.spec.vcpus,
req.spec.restart,
writable_root,
);
let mounts = self.stage_volumes(&req.spec).await?;
if !mounts.is_empty() {
host_config.mounts = Some(mounts);
}
let mut config = Config {
image: Some(reference),
cmd: Some(req.spec.entrypoint.clone()),
env: Some(env),
..Default::default()
};
if self.endpoint == DockerEndpoint::Published {
config.exposed_ports = Some(HashMap::from([(port_key.clone(), HashMap::new())]));
host_config.port_bindings = Some(HashMap::from([(
port_key.clone(),
Some(vec![PortBinding {
host_ip: Some("127.0.0.1".to_string()),
host_port: Some("0".to_string()),
}]),
)]));
}
config.host_config = Some(host_config);
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 (host, endpoint_port) = match self.endpoint {
DockerEndpoint::Published => (
"127.0.0.1".to_string(),
self.published_host_port(&created.id, &port_key).await?,
),
DockerEndpoint::Bridge => (self.container_ip(&created.id).await?, port),
};
Ok(Instance {
handle: InstanceHandle {
workload: req.workload.clone(),
replica: req.replica,
backend_ref: encode_ref(&name, &host, endpoint_port),
},
endpoint: Endpoint {
scheme: Scheme::Http,
host,
port: endpoint_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()))
}
async fn published_host_port(&self, id: &str, port_key: &str) -> Result<u16, BackendError> {
let info = self
.docker
.inspect_container(id, None)
.await
.map_err(|e| BackendError::Launch(format!("inspect {id}: {e}")))?;
info.network_settings
.and_then(|ns| ns.ports)
.and_then(|mut ports| ports.remove(port_key).flatten())
.and_then(|bindings| bindings.into_iter().next())
.and_then(|b| b.host_port)
.and_then(|hp| hp.parse::<u16>().ok())
.ok_or_else(|| {
BackendError::Launch(format!("no published host port for {port_key} on {id}"))
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn docker_endpoint_defaults_to_published_and_parses_lowercase() {
assert_eq!(DockerEndpoint::default(), DockerEndpoint::Published);
assert_eq!(
serde_json::from_str::<DockerEndpoint>("\"published\"").unwrap(),
DockerEndpoint::Published
);
assert_eq!(
serde_json::from_str::<DockerEndpoint>("\"bridge\"").unwrap(),
DockerEndpoint::Bridge
);
}
#[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, false);
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, false).nano_cpus,
Some(1_000_000_000)
);
}
#[test]
fn writable_root_relaxes_only_the_read_only_root() {
let hc = hardened_host_config(256, 2, RestartPolicy::Never, true);
assert_eq!(hc.readonly_rootfs, Some(false));
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.pids_limit, Some(MAX_PIDS));
}
#[test]
fn writable_root_is_off_by_default_on_the_backend() {
let docker = Docker::connect_with_defaults().unwrap();
let backend = DockerBackend::with_client(docker);
assert!(!backend.writable_root_allowed);
assert!(
backend
.with_writable_root_allowed(true)
.writable_root_allowed,
"the single-tenant posture opts in"
);
}
#[test]
fn named_volume_mode_builds_a_prefixed_daemon_volume_mount() {
let vol = VolumeRef {
name: "db".into(),
mount: "/data".into(),
size_mib: 64,
};
let m = volume_mount(&vol, DockerVolumeMode::Named, Path::new("/srv/data"));
assert_eq!(m.typ, Some(MountTypeEnum::VOLUME));
assert_eq!(m.source.as_deref(), Some("boatramp-db"));
assert_eq!(m.target.as_deref(), Some("/data"));
assert_eq!(m.read_only, Some(false), "a persistent volume is writable");
}
#[test]
fn bind_volume_mode_builds_a_host_path_mount() {
let vol = VolumeRef {
name: "db".into(),
mount: "/data".into(),
size_mib: 64,
};
let m = volume_mount(&vol, DockerVolumeMode::Bind, Path::new("/srv/data"));
assert_eq!(m.typ, Some(MountTypeEnum::BIND));
assert_eq!(m.source.as_deref(), Some("/srv/data/compute/volumes/db"));
assert_eq!(m.target.as_deref(), Some("/data"));
assert_eq!(m.read_only, Some(false));
}
#[test]
fn validate_volume_rejects_traversal_in_name_and_mount() {
assert!(validate_volume("db", "/data").is_ok());
assert!(validate_volume("cache-1", "/var/lib/app").is_ok());
assert!(validate_volume("../etc", "/data").is_err());
assert!(validate_volume("a/b", "/data").is_err());
assert!(validate_volume("db", "relative").is_err());
assert!(validate_volume("db", "/data/../etc").is_err());
}
#[test]
fn volume_mode_defaults_to_named_and_parses_lowercase() {
assert_eq!(DockerVolumeMode::default(), DockerVolumeMode::Named);
assert_eq!(
serde_json::from_str::<DockerVolumeMode>("\"named\"").unwrap(),
DockerVolumeMode::Named
);
assert_eq!(
serde_json::from_str::<DockerVolumeMode>("\"bind\"").unwrap(),
DockerVolumeMode::Bind
);
}
#[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)
);
}
}