use std::sync::Arc;
#[cfg(target_os = "linux")]
struct PostureKernelVerifier {
strict: bool,
signing_keys: Vec<String>,
allowed_hashes: Vec<String>,
daemon: Option<Arc<boatramp_server::DaemonRuntime>>,
}
#[cfg(target_os = "linux")]
impl std::fmt::Debug for PostureKernelVerifier {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PostureKernelVerifier")
.field("strict", &self.strict)
.field("signing_keys", &self.signing_keys.len())
.field("allowed_hashes", &self.allowed_hashes.len())
.field("has_daemon", &self.daemon.is_some())
.finish()
}
}
#[cfg(target_os = "linux")]
impl boatramp_firecracker::KernelVerifier for PostureKernelVerifier {
fn verify(&self, bytes: &[u8], expected_hash: &str) -> std::result::Result<(), String> {
let sig = self
.daemon
.as_ref()
.and_then(|d| d.effective().default_kernel.clone())
.filter(|dk| dk.sha256 == expected_hash)
.and_then(|dk| dk.sig);
let kref = boatramp_core::daemon_config::KernelRef {
source: expected_hash.to_string(),
sha256: expected_hash.to_string(),
sig,
};
boatramp_core::kernel_trust::verify_kernel(
bytes,
&kref,
self.strict,
&self.signing_keys,
&self.allowed_hashes,
)
.map_err(|e| e.to_string())
}
}
#[cfg(target_os = "macos")]
struct VzPostureKernelVerifier {
strict: bool,
signing_keys: Vec<String>,
allowed_hashes: Vec<String>,
daemon: Option<Arc<boatramp_server::DaemonRuntime>>,
}
#[cfg(target_os = "macos")]
impl std::fmt::Debug for VzPostureKernelVerifier {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("VzPostureKernelVerifier")
.field("strict", &self.strict)
.field("signing_keys", &self.signing_keys.len())
.field("allowed_hashes", &self.allowed_hashes.len())
.field("has_daemon", &self.daemon.is_some())
.finish()
}
}
#[cfg(target_os = "macos")]
impl boatramp_vz::KernelVerifier for VzPostureKernelVerifier {
fn verify(&self, bytes: &[u8], expected_hash: &str) -> std::result::Result<(), String> {
let sig = self
.daemon
.as_ref()
.and_then(|d| d.effective().default_kernel.clone())
.filter(|dk| dk.sha256 == expected_hash)
.and_then(|dk| dk.sig);
let kref = boatramp_core::daemon_config::KernelRef {
source: expected_hash.to_string(),
sha256: expected_hash.to_string(),
sig,
};
boatramp_core::kernel_trust::verify_kernel(
bytes,
&kref,
self.strict,
&self.signing_keys,
&self.allowed_hashes,
)
.map_err(|e| e.to_string())
}
}
#[cfg(target_os = "macos")]
fn macos_supports_vz() -> bool {
if cfg!(not(target_arch = "aarch64")) {
return false;
}
let major = sysctl_string("kern.osproductversion")
.and_then(|v| v.split('.').next().and_then(|m| m.parse::<u32>().ok()));
matches!(major, Some(m) if m >= 15)
}
#[cfg(target_os = "macos")]
fn sysctl_string(name: &str) -> Option<String> {
let out = std::process::Command::new("sysctl")
.args(["-n", name])
.output()
.ok()?;
if !out.status.success() {
return None;
}
Some(String::from_utf8_lossy(&out.stdout).trim().to_string())
}
pub async fn build_compute(
cfg: Option<&crate::config::ComputeConfig>,
storage: std::sync::Arc<dyn boatramp_core::Storage>,
data_dir: &std::path::Path,
node_id: u64,
strict: bool,
daemon: Option<Arc<boatramp_server::DaemonRuntime>>,
worker_exe: Option<&std::path::Path>,
) -> (
boatramp_core::compute::BackendRegistry,
boatramp_core::compute::Node,
) {
use boatramp_core::compute::{BackendKind, BackendRegistry, Node};
let mut backends: BackendRegistry = std::collections::BTreeMap::new();
let empty_node = |id| Node {
id,
region: None,
labels: std::collections::BTreeMap::new(),
free_vcpus: 0,
free_mem_mib: 0,
backends: Vec::new(),
};
let Some(cfg) = cfg else {
return (backends, empty_node(node_id));
};
match boatramp_docker::DockerBackend::connect() {
Ok(docker) => {
let docker = docker
.with_endpoint(cfg.docker_endpoint)
.with_volume_mode(cfg.docker_volume_mode)
.with_data_dir(data_dir)
.with_writable_root_allowed(!strict)
.with_cap_add_allowed(!strict);
if docker.reachable().await {
backends.insert("docker".to_string(), std::sync::Arc::new(docker));
} else {
tracing::debug!("no reachable docker daemon; skipping docker backend");
}
}
Err(e) => tracing::debug!(%e, "docker backend unavailable"),
}
#[cfg(target_os = "linux")]
let shared_ip_authority: Option<boatramp_core::ipam::IpAuthority> =
match boatramp_core::ipam::IpAuthority::new(&cfg.subnet) {
Ok(a) => Some(a),
Err(e) => {
tracing::warn!(%e, subnet = %cfg.subnet, "bad compute subnet; container + embedded-VMM backends disabled");
None
}
};
#[cfg(target_os = "linux")]
let bridge_ready = match &shared_ip_authority {
Some(authority) => {
match boatramp_container::ensure_bridge(
&cfg.bridge,
authority.gateway(),
authority.prefix_len(),
)
.await
{
Ok(()) => true,
Err(e) => {
tracing::warn!(%e, bridge = %cfg.bridge, "could not create the compute bridge (need CAP_NET_ADMIN); container + embedded-VMM backends disabled");
false
}
}
}
None => false,
};
#[cfg(target_os = "linux")]
if bridge_ready {
match worker_exe.map_or_else(std::env::current_exe, |p| Ok(p.to_path_buf())) {
Ok(self_exe) => match boatramp_container::ContainerBackend::new(
storage.clone(),
data_dir.to_path_buf(),
cfg.bridge.clone(),
&cfg.subnet,
self_exe,
) {
Ok(c) => {
let c = c.with_cap_add_allowed(!strict);
let c = c.with_internal_dns(cfg.internal_dns.then(|| cfg.dns_domain.clone()));
let c = match &shared_ip_authority {
Some(a) => c.with_ip_authority(a.clone()),
None => c,
};
backends.insert("container".to_string(), std::sync::Arc::new(c));
}
Err(e) => tracing::warn!(%e, "container backend unavailable"),
},
Err(e) => tracing::warn!(%e, "current_exe for container backend"),
}
}
#[cfg(all(target_os = "linux", target_arch = "x86_64"))]
if bridge_ready && std::path::Path::new("/dev/kvm").exists() {
match (
worker_exe.map_or_else(std::env::current_exe, |p| Ok(p.to_path_buf())),
boatramp_core::ipam::IpPool::new(&cfg.subnet),
) {
(Ok(self_exe), Ok(pool)) => {
let gateway = pool.gateway().to_string();
let verifier: Arc<dyn boatramp_firecracker::KernelVerifier> =
Arc::new(PostureKernelVerifier {
strict,
signing_keys: cfg.kernel_signing_pubkeys.clone(),
allowed_hashes: cfg.kernel_allowed_hashes.clone(),
daemon: daemon.clone(),
});
match boatramp_firecracker::EmbeddedVmmBackend::new(
storage.clone(),
self_exe, data_dir.to_path_buf(),
cfg.bridge.clone(),
gateway,
&cfg.subnet,
verifier,
) {
Ok(vmm) => {
let vmm = match &shared_ip_authority {
Some(a) => vmm.with_ip_authority(a.clone()),
None => vmm,
};
backends.insert("vmm-embedded".to_string(), std::sync::Arc::new(vmm));
}
Err(e) => tracing::warn!(%e, "embedded VMM backend unavailable"),
}
}
(Err(e), _) => tracing::warn!(%e, "current_exe for VMM backend"),
(_, Err(e)) => tracing::warn!(%e, "bad compute subnet for VMM backend"),
}
} else {
tracing::debug!("no /dev/kvm; skipping embedded VMM backend");
}
#[cfg(target_os = "macos")]
if macos_supports_vz() {
match (
worker_exe.map_or_else(std::env::current_exe, |p| Ok(p.to_path_buf())),
boatramp_core::ipam::IpPool::new(&cfg.subnet),
) {
(Ok(self_exe), Ok(_pool)) => {
let verifier: Arc<dyn boatramp_vz::KernelVerifier> =
Arc::new(VzPostureKernelVerifier {
strict,
signing_keys: cfg.kernel_signing_pubkeys.clone(),
allowed_hashes: cfg.kernel_allowed_hashes.clone(),
daemon: daemon.clone(),
});
match boatramp_vz::VzBackend::new(
storage.clone(),
self_exe, data_dir.to_path_buf(),
&cfg.subnet, verifier,
) {
Ok(vz) => {
let vz = vz.with_writable_root_allowed(!strict);
backends.insert("vmm-vz".to_string(), std::sync::Arc::new(vz));
}
Err(e) => tracing::warn!(%e, "macOS VMM backend unavailable"),
}
}
(Err(e), _) => tracing::warn!(%e, "current_exe for macOS VMM backend"),
(_, Err(e)) => tracing::warn!(%e, "bad compute subnet for macOS VMM backend"),
}
} else {
tracing::debug!("not Apple silicon + macOS 15+; skipping macOS VMM backend");
}
let _ = (&storage, data_dir); #[cfg(not(any(all(target_os = "linux", target_arch = "x86_64"), target_os = "macos")))]
let _ = (strict, &daemon);
let free_vcpus = if cfg.vcpus > 0 {
cfg.vcpus
} else {
std::thread::available_parallelism()
.map(|n| n.get() as u32)
.unwrap_or(1)
};
let free_mem_mib = if cfg.mem_mib > 0 { cfg.mem_mib } else { 1024 };
let advertised: Vec<BackendKind> = backends
.iter()
.map(|(id, b)| {
let caps = b.capabilities();
BackendKind {
id: id.clone(),
isolation: caps.isolation,
persistent_volumes: caps.persistent_volumes,
scale_to_zero: caps.scale_to_zero,
}
})
.collect();
tracing::info!(backends = ?advertised, free_vcpus, free_mem_mib, "compute node inventory");
let node = Node {
id: node_id,
region: cfg.region.clone(),
labels: std::collections::BTreeMap::new(),
free_vcpus,
free_mem_mib,
backends: advertised,
};
(backends, node)
}
pub async fn adopt_running_replica_ips(
deploy: &boatramp_core::deploy::DeployStore,
backends: &boatramp_core::compute::BackendRegistry,
) {
use std::collections::BTreeMap;
use std::net::Ipv4Addr;
let states = match deploy.list_all_replica_states().await {
Ok(s) => s,
Err(e) => {
tracing::warn!(%e, "could not read replica states for IP adoption; \
the reconcile loop starts without adopting in-use IPs");
return;
}
};
let mut by_backend: BTreeMap<String, Vec<(String, String, u32, Ipv4Addr)>> = BTreeMap::new();
for st in &states {
let ip = st.endpoint.host.parse::<Ipv4Addr>().ok().or_else(|| {
st.handle
.backend_ref
.split(':')
.next()
.and_then(|s| s.parse::<Ipv4Addr>().ok())
});
if let Some(ip) = ip {
by_backend.entry(st.backend.clone()).or_default().push((
st.handle.project.clone(),
st.handle.workload.clone(),
st.handle.replica,
ip,
));
}
}
for (backend_id, replicas) in by_backend {
if let Some(backend) = backends.get(&backend_id) {
backend.reserve_in_use(&replicas).await;
tracing::info!(
backend = %backend_id,
count = replicas.len(),
"adopted in-use compute IPs into the backend pool"
);
}
}
}
#[cfg(target_os = "linux")]
pub fn spawn_internal_dns(
cfg: Option<&crate::config::ComputeConfig>,
backends: &boatramp_core::compute::BackendRegistry,
deploy: &boatramp_core::deploy::DeployStore,
) -> Option<tokio::task::JoinHandle<()>> {
let cfg = cfg?;
if !cfg.internal_dns || !backends.contains_key("container") {
return None;
}
let gateway = match boatramp_core::ipam::IpPool::new(&cfg.subnet) {
Ok(pool) => pool.gateway(),
Err(e) => {
tracing::warn!(%e, subnet = %cfg.subnet, "internal DNS: bad compute subnet; resolver not started");
return None;
}
};
let upstream: std::net::SocketAddr = match cfg.dns_upstream.parse() {
Ok(a) => a,
Err(e) => {
tracing::warn!(%e, upstream = %cfg.dns_upstream, "internal DNS: bad dns_upstream (want host:port); resolver not started");
return None;
}
};
let source: std::sync::Arc<dyn boatramp_container::dns_server::InternalDnsSource> =
std::sync::Arc::new(DeployDnsSource::new(deploy.clone()));
let domain = cfg.dns_domain.clone();
Some(tokio::spawn(async move {
if let Err(e) =
boatramp_container::dns_server::serve(gateway, upstream, domain, source).await
{
tracing::warn!(%e, "internal DNS resolver exited (bind/setup error); \
guests keep their static resolv.conf peers");
}
}))
}
#[cfg(not(target_os = "linux"))]
pub fn spawn_internal_dns(
_cfg: Option<&crate::config::ComputeConfig>,
_backends: &boatramp_core::compute::BackendRegistry,
_deploy: &boatramp_core::deploy::DeployStore,
) -> Option<tokio::task::JoinHandle<()>> {
None
}
#[cfg(target_os = "linux")]
pub struct DeployDnsSource {
deploy: boatramp_core::deploy::DeployStore,
}
#[cfg(target_os = "linux")]
impl DeployDnsSource {
pub fn new(deploy: boatramp_core::deploy::DeployStore) -> Self {
Self { deploy }
}
}
#[cfg(target_os = "linux")]
#[async_trait::async_trait]
impl boatramp_container::dns_server::InternalDnsSource for DeployDnsSource {
async fn snapshot(&self) -> boatramp_container::dns_server::DnsFleet {
use boatramp_container::dns::ResolvedAddrs;
use boatramp_container::dns_server::DnsFleet;
use boatramp_core::compute::ReplicaPhase;
use std::net::Ipv4Addr;
let mut fleet = DnsFleet::default();
let states = match self.deploy.list_all_replica_states().await {
Ok(s) => s,
Err(e) => {
tracing::warn!(%e, "internal DNS: could not read replica states; \
answering forward-only this query");
return fleet;
}
};
for st in &states {
let v4 = st.endpoint.host.parse::<Ipv4Addr>().ok().or_else(|| {
st.handle
.backend_ref
.split(':')
.next()
.and_then(|s| s.parse::<Ipv4Addr>().ok())
});
let Some(ip) = v4 else { continue };
let key = (st.handle.project.clone(), st.handle.workload.clone());
fleet.owners.insert(ip, key.clone());
if st.phase == ReplicaPhase::Running && st.healthy {
fleet
.addrs
.entry(key)
.or_insert_with(ResolvedAddrs::default)
.v4
.push(ip);
}
}
fleet
}
}
pub struct NodeComputeExec {
backends: boatramp_core::compute::BackendRegistry,
deploy: boatramp_core::deploy::DeployStore,
}
impl NodeComputeExec {
pub fn new(
backends: boatramp_core::compute::BackendRegistry,
deploy: boatramp_core::deploy::DeployStore,
) -> Self {
Self { backends, deploy }
}
}
#[async_trait::async_trait]
impl boatramp_core::compute::ComputeExec for NodeComputeExec {
async fn exec(
&self,
project: &str,
workload: &str,
argv: &[String],
stdin: Option<&[u8]>,
) -> Result<boatramp_core::compute::ExecOutput, boatramp_core::compute::ExecError> {
use boatramp_core::compute::{BackendError, ExecError, ReplicaPhase};
use boatramp_core::project::ProjectRef;
let states = self
.deploy
.list_replica_states(ProjectRef::new(project), workload)
.await
.map_err(|e| ExecError::Other(e.to_string()))?;
let target = states
.iter()
.find(|s| s.phase == ReplicaPhase::Running && s.healthy)
.or_else(|| states.iter().find(|s| s.phase == ReplicaPhase::Running))
.ok_or_else(|| ExecError::NoReplica(workload.to_string()))?;
let backend = self
.backends
.get(&target.backend)
.ok_or_else(|| ExecError::Unsupported(target.backend.clone()))?;
match backend.exec(&target.handle, argv, stdin).await {
Ok(out) => Ok(out),
Err(BackendError::Unsupported) => Err(ExecError::Unsupported(target.backend.clone())),
Err(e) => Err(ExecError::Other(e.to_string())),
}
}
}
pub struct NodeComputeVolumes {
backends: boatramp_core::compute::BackendRegistry,
deploy: boatramp_core::deploy::DeployStore,
}
impl NodeComputeVolumes {
pub fn new(
backends: boatramp_core::compute::BackendRegistry,
deploy: boatramp_core::deploy::DeployStore,
) -> Self {
Self { backends, deploy }
}
async fn referenced_volume_names(
&self,
) -> Result<std::collections::BTreeSet<String>, boatramp_core::compute::VolumeError> {
use boatramp_core::compute::VolumeError;
let mut names = std::collections::BTreeSet::new();
let workloads = self
.deploy
.list_compute_workloads_all()
.await
.map_err(|e| VolumeError::Other(e.to_string()))?;
for (_project, workload) in workloads {
let spec = self
.deploy
.get_compute_spec(&workload.active)
.await
.map_err(|e| VolumeError::Other(e.to_string()))?;
if let Some(spec) = spec {
for vol in spec.volumes {
names.insert(vol.name);
}
}
}
Ok(names)
}
}
#[async_trait::async_trait]
impl boatramp_core::compute::ComputeVolumes for NodeComputeVolumes {
async fn list(
&self,
) -> Result<Vec<boatramp_core::compute::VolumeStatus>, boatramp_core::compute::VolumeError>
{
use boatramp_core::compute::{VolumeError, VolumeStatus};
let referenced = self.referenced_volume_names().await?;
let mut by_name: std::collections::BTreeMap<String, u64> =
std::collections::BTreeMap::new();
for backend in self.backends.values() {
let vols = backend
.list_volumes()
.await
.map_err(|e| VolumeError::Other(e.to_string()))?;
for v in vols {
let slot = by_name.entry(v.name).or_insert(0);
*slot = (*slot).max(v.size_bytes);
}
}
Ok(by_name
.into_iter()
.map(|(name, size_bytes)| VolumeStatus {
in_use: referenced.contains(&name),
info: boatramp_core::compute::VolumeInfo { name, size_bytes },
})
.collect())
}
async fn remove(
&self,
name: &str,
force: bool,
) -> Result<bool, boatramp_core::compute::VolumeError> {
use boatramp_core::compute::{BackendError, VolumeError};
if !force && self.referenced_volume_names().await?.contains(name) {
return Err(VolumeError::InUse(name.to_string()));
}
let mut existed = false;
let mut any_supported = false;
for backend in self.backends.values() {
match backend.remove_volume(name).await {
Ok(removed) => {
any_supported = true;
existed |= removed;
}
Err(BackendError::Unsupported) => {}
Err(e) => return Err(VolumeError::Other(e.to_string())),
}
}
if !any_supported {
return Err(VolumeError::Unsupported);
}
Ok(existed)
}
}
#[cfg(test)]
mod tests {
use super::*;
use async_trait::async_trait;
use boatramp_core::compute::{
Artifact, BackendError, Capabilities, ComputeBackend, ComputeSpec, ComputeVolumes,
ComputeWorkload, Health, Instance, InstanceHandle, IsolationClass, IsolationRequirement,
LaunchRequest, RestartPolicy, RootSource, VolumeError, VolumeInfo, VolumeRef,
};
use boatramp_core::deploy::DeployStore;
use boatramp_core::project::ProjectRef;
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
struct NullStorage;
#[async_trait]
impl Storage for NullStorage {
async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn get_range(
&self,
_: &str,
_: u64,
_: Option<u64>,
) -> Result<GetObject, StorageError> {
Err(StorageError::unsupported("range"))
}
async fn put(
&self,
_: &str,
_: ByteStream,
_: PutMeta,
) -> Result<ObjectMeta, StorageError> {
Err(StorageError::unsupported("put"))
}
async fn head(&self, _: &str) -> Result<ObjectMeta, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn delete(&self, _: &str) -> Result<(), StorageError> {
Ok(())
}
async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
Ok(Vec::new())
}
}
struct FakeVolumeBackend {
vols: Mutex<BTreeMap<String, u64>>,
}
impl FakeVolumeBackend {
fn with(names: &[(&str, u64)]) -> Self {
Self {
vols: Mutex::new(names.iter().map(|(n, s)| (n.to_string(), *s)).collect()),
}
}
}
#[async_trait]
impl ComputeBackend for FakeVolumeBackend {
fn id(&self) -> &'static str {
"container"
}
fn capabilities(&self) -> Capabilities {
Capabilities {
isolation: IsolationClass::Namespace,
scale_to_zero: false,
persistent_volumes: true,
max_vcpus: None,
max_mem_mib: None,
}
}
async fn materialize(&self, _: &ComputeSpec) -> Result<Artifact, BackendError> {
Err(BackendError::Unsupported)
}
async fn launch(&self, _: &LaunchRequest) -> Result<Instance, BackendError> {
Err(BackendError::Unsupported)
}
async fn stop(&self, _: &InstanceHandle) -> Result<(), BackendError> {
Ok(())
}
async fn health(&self, _: &InstanceHandle) -> Result<Health, BackendError> {
Ok(Health::Unknown)
}
async fn list_volumes(&self) -> Result<Vec<VolumeInfo>, BackendError> {
Ok(self
.vols
.lock()
.unwrap()
.iter()
.map(|(name, size)| VolumeInfo {
name: name.clone(),
size_bytes: *size,
})
.collect())
}
async fn remove_volume(&self, name: &str) -> Result<bool, BackendError> {
Ok(self.vols.lock().unwrap().remove(name).is_some())
}
}
fn spec_with_volume(vol: Option<&str>) -> ComputeSpec {
ComputeSpec {
version: 1,
root: RootSource::Image("img".into()),
kernel: String::new(),
kernel_cmdline: None,
vcpus: 1,
mem_mib: 64,
entrypoint: vec![],
env: BTreeMap::new(),
port: 8080,
restart: RestartPolicy::Always,
startup_grace_secs: 30,
scale_to_zero: false,
volumes: vol
.map(|n| {
vec![VolumeRef {
mount: "/data".into(),
name: n.into(),
size_mib: 128,
}]
})
.unwrap_or_default(),
writable_root: false,
cap_add: vec![],
user: None,
isolation: IsolationRequirement::Trusted,
prefer_backend: None,
bindings: vec![],
}
}
async fn setup(referenced: Option<&str>, backend_vols: &[(&str, u64)]) -> NodeComputeVolumes {
let store = DeployStore::new(
Arc::new(NullStorage),
Arc::new(boatramp_core::kv::MemoryKv::new()),
);
let spec = spec_with_volume(referenced);
let hash = store.put_compute_spec(&spec).await.expect("put spec");
let workload = ComputeWorkload {
version: 1,
name: "wl".into(),
active: hash,
replicas: 1,
placement: Default::default(),
};
store
.set_compute_workload(ProjectRef::DEFAULT, &workload)
.await
.expect("set workload");
let mut backends: boatramp_core::compute::BackendRegistry = BTreeMap::new();
backends.insert(
"container".into(),
Arc::new(FakeVolumeBackend::with(backend_vols)) as Arc<dyn ComputeBackend>,
);
NodeComputeVolumes::new(backends, store)
}
#[tokio::test]
async fn list_flags_referenced_volume_in_use_and_orphan_free() {
let vols = setup(Some("data"), &[("data", 100), ("old", 50)]).await;
let listed = vols.list().await.expect("list");
assert_eq!(listed.len(), 2);
let data = listed.iter().find(|v| v.info.name == "data").unwrap();
let old = listed.iter().find(|v| v.info.name == "old").unwrap();
assert!(data.in_use, "spec-referenced volume is in use");
assert_eq!(data.info.size_bytes, 100);
assert!(!old.in_use, "unreferenced volume is orphaned");
assert_eq!(old.info.size_bytes, 50);
}
#[tokio::test]
async fn remove_refuses_in_use_without_force_and_allows_with_force() {
let vols = setup(Some("data"), &[("data", 100)]).await;
assert!(matches!(
vols.remove("data", false).await,
Err(VolumeError::InUse(n)) if n == "data"
));
assert!(vols
.list()
.await
.unwrap()
.iter()
.any(|v| v.info.name == "data"));
assert!(vols.remove("data", true).await.expect("forced remove"));
assert!(vols.list().await.unwrap().is_empty());
}
#[tokio::test]
async fn remove_orphan_succeeds_and_absent_reports_false() {
let vols = setup(None, &[("old", 50)]).await;
assert!(vols.remove("old", false).await.expect("remove orphan"));
assert!(!vols.remove("gone", false).await.expect("remove absent"));
}
struct AdoptSpyBackend {
adopted: Mutex<Vec<(String, String, u32, std::net::Ipv4Addr)>>,
}
#[async_trait]
impl ComputeBackend for AdoptSpyBackend {
fn id(&self) -> &'static str {
"container"
}
fn capabilities(&self) -> Capabilities {
Capabilities {
isolation: IsolationClass::Namespace,
scale_to_zero: false,
persistent_volumes: true,
max_vcpus: None,
max_mem_mib: None,
}
}
async fn materialize(&self, _: &ComputeSpec) -> Result<Artifact, BackendError> {
Err(BackendError::Unsupported)
}
async fn reserve_in_use(&self, replicas: &[(String, String, u32, std::net::Ipv4Addr)]) {
self.adopted.lock().unwrap().extend_from_slice(replicas);
}
async fn launch(&self, _: &LaunchRequest) -> Result<Instance, BackendError> {
Err(BackendError::Unsupported)
}
async fn stop(&self, _: &InstanceHandle) -> Result<(), BackendError> {
Ok(())
}
async fn health(&self, _: &InstanceHandle) -> Result<Health, BackendError> {
Ok(Health::Unknown)
}
}
#[tokio::test]
async fn adopt_running_replica_ips_feeds_each_backends_in_use_addresses() {
use boatramp_core::compute::{Endpoint, ObservedInstance, ReplicaPhase, Scheme};
use std::net::Ipv4Addr;
let store = DeployStore::new(
Arc::new(NullStorage),
Arc::new(boatramp_core::kv::MemoryKv::new()),
);
let mk = |proj: &str, wl: &str, rep: u32, backend: &str, ip: &str, phase: ReplicaPhase| {
(
proj.to_string(),
ObservedInstance {
handle: InstanceHandle {
project: proj.into(),
workload: wl.into(),
replica: rep,
backend_ref: format!("{ip}:5432"),
},
node: 1,
backend: backend.into(),
endpoint: Endpoint {
scheme: Scheme::Http,
host: ip.into(),
port: 5432,
},
region: None,
healthy: true,
started_at: None,
phase,
snapshot: None,
},
)
};
for (proj, st) in [
mk(
"default",
"pg-a",
0,
"container",
"10.0.0.2",
ReplicaPhase::Running,
),
mk(
"acme",
"web",
0,
"container",
"10.0.0.3",
ReplicaPhase::Zero,
), mk(
"default",
"vm",
0,
"vmm-embedded",
"10.0.0.9",
ReplicaPhase::Running,
),
] {
store
.set_replica_state(ProjectRef::new(&proj), &st)
.await
.expect("persist replica state");
}
let container = Arc::new(AdoptSpyBackend {
adopted: Mutex::new(Vec::new()),
});
let mut backends: boatramp_core::compute::BackendRegistry = BTreeMap::new();
backends.insert(
"container".into(),
container.clone() as Arc<dyn ComputeBackend>,
);
adopt_running_replica_ips(&store, &backends).await;
let got = container.adopted.lock().unwrap().clone();
assert!(got.contains(&(
"default".into(),
"pg-a".into(),
0,
Ipv4Addr::new(10, 0, 0, 2)
)));
assert!(got.contains(&("acme".into(), "web".into(), 0, Ipv4Addr::new(10, 0, 0, 3))));
assert!(
!got.iter()
.any(|(_, _, _, ip)| *ip == Ipv4Addr::new(10, 0, 0, 9)),
"another backend's replica IP must not be adopted by the container backend"
);
assert_eq!(got.len(), 2);
}
}