mod convert;
mod events;
mod execution;
mod files;
mod snapshots;
mod template;
mod templates;
mod window;
use std::collections::HashMap;
use std::sync::{Arc, Mutex, Weak};
use arcbox_computer_runtime::agent::VmProtoAgentFactory;
use arcbox_computer_runtime::{
NodeEnvironment, RootfsBuilder, RootfsPaths, SandboxManager, SandboxMountSpec,
SandboxNetworkSpec, SandboxSpec, SandboxState, VmmError,
};
use arcbox_connect::sandbox_v1;
use arcbox_fc_driver::{FcDriver, FcDriverConfig};
use arcbox_snapshot::snapshot_cow::{BlockTools, BusyboxBlockTools, CowManager, CowOptions};
use arcbox_tap_net::{IptablesLegacy, TapNetwork};
use buffa::Message;
use tokio::sync::{Mutex as AsyncMutex, OwnedMutexGuard};
use crate::config::GuestConfig;
use crate::create_registry::{CreateRegistry, Reserve as CreateReserve};
use crate::error::SandboxError;
const STARTUP_CLEANUP_ID: &str = "$startup";
pub fn probe_kvm() -> Result<(), String> {
match std::fs::OpenOptions::new()
.read(true)
.write(true)
.open("/dev/kvm")
{
Ok(_) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Err(
"sandboxes require nested virtualization (/dev/kvm is missing in the guest): \
use the VZ backend on Apple Silicon M3 or newer with macOS 15+; \
the HV backend and Intel/M1/M2 hosts cannot run sandboxes"
.into(),
),
Err(e) => Err(format!(
"/dev/kvm exists but cannot be opened ({e}); sandboxes are unavailable"
)),
}
}
pub struct SandboxService {
manager: Arc<SandboxManager>,
creates: Arc<CreateRegistry>,
operations: SandboxOperationLocks,
template_builds: Mutex<std::collections::HashSet<String>>,
default_rootfs: String,
rootfs: RootfsBuilder,
}
const VM_AGENT_BIN: &str = "/arcbox/bin/vm-agent";
const GUEST_BUSYBOX: &str = "/bin/busybox";
const ROOTFS_CACHE_DIR: &str = "/var/lib/arcbox/sandbox";
pub fn block_tools() -> Arc<dyn BlockTools> {
Arc::new(BusyboxBlockTools::default())
}
pub fn rootfs_builder(block_tools: Arc<dyn BlockTools>) -> RootfsBuilder {
RootfsBuilder::new(
RootfsPaths {
vm_agent: VM_AGENT_BIN.into(),
cache_dir: ROOTFS_CACHE_DIR.into(),
busybox: GUEST_BUSYBOX.into(),
},
block_tools,
)
}
pub fn node_environment(
config: &GuestConfig,
block_tools: Arc<dyn BlockTools>,
) -> anyhow::Result<NodeEnvironment> {
let runtime = &config.runtime;
let data_dir = std::path::Path::new(&runtime.firecracker.data_dir);
let network = TapNetwork::with_quarantine_dir(
&runtime.network.cidr,
&runtime.network.gateway,
runtime.network.dns.clone(),
data_dir.join("sandbox-network-quarantine"),
config.adapters.sandbox_datapath,
Arc::new(IptablesLegacy::default()),
)?;
let mut cow_options = CowOptions::new(data_dir);
cow_options.block_tools = block_tools;
if let Some(candidates) = &runtime.firecracker.dmsetup_candidates {
cow_options.dmsetup_candidates = candidates.iter().map(std::path::PathBuf::from).collect();
}
Ok(NodeEnvironment {
driver: Arc::new(FcDriver::new(FcDriverConfig::from(&config.adapters))),
network: Arc::new(network),
agent: Arc::new(VmProtoAgentFactory::default()),
cow_manager: Arc::new(CowManager::new(cow_options)?),
})
}
#[derive(Default)]
struct SandboxOperationLocks {
entries: Mutex<HashMap<String, Weak<AsyncMutex<()>>>>,
}
impl SandboxOperationLocks {
async fn lock(&self, id: &str) -> Option<OwnedMutexGuard<()>> {
if id.is_empty() {
return None;
}
let lock = {
let mut entries = self.entries.lock().unwrap();
entries.retain(|_, lock| lock.strong_count() > 0);
if let Some(lock) = entries.get(id).and_then(Weak::upgrade) {
lock
} else {
let lock = Arc::new(AsyncMutex::new(()));
entries.insert(id.to_owned(), Arc::downgrade(&lock));
lock
}
};
Some(lock.lock_owned().await)
}
}
impl SandboxService {
pub async fn lock_operation(&self, id: &str) -> Option<OwnedMutexGuard<()>> {
self.operations.lock(id).await
}
pub fn new(config: GuestConfig) -> anyhow::Result<Self> {
let default_rootfs = config.runtime.defaults.rootfs.clone();
let block_tools = block_tools();
let rootfs = rootfs_builder(Arc::clone(&block_tools));
let environment = node_environment(&config, block_tools)?;
let manager = SandboxManager::new(config.runtime, environment)
.map_err(|e| anyhow::anyhow!("{e}"))?
.into_shared();
let creates = Arc::new(CreateRegistry::default());
Ok(Self {
manager,
creates,
operations: SandboxOperationLocks::default(),
template_builds: Mutex::new(std::collections::HashSet::new()),
default_rootfs,
rootfs,
})
}
pub async fn create(
&self,
payload: &[u8],
) -> Result<sandbox_v1::CreateSandboxResponse, SandboxError> {
let mut request = sandbox_v1::CreateSandboxRequest::decode_from_slice(payload)
.map_err(|e| SandboxError::Decode(e.to_string()))?;
if create_uses_network(&request) {
self.manager.wait_startup_cleanup_complete().await;
}
let source = self.resolve_template_source(&mut request)?;
let _operation = self.operations.lock(&request.id).await;
let create_key = crate::create_key::create_key(&request);
if request.id.is_empty() {
return self.create_once(request, source, &create_key).await;
}
let id = request.id.clone();
loop {
self.clear_stale_completed_create(&id);
match self.creates.reserve(&request).map_err(|_| {
SandboxError::AlreadyExists(format!(
"sandbox '{id}' (id reused for a different create request)"
))
})? {
CreateReserve::Existing(response) => return Ok(response),
CreateReserve::Slot(slot) => {
let response = self.create_once(request, source, &create_key).await?;
slot.commit(&response);
return Ok(response);
}
CreateReserve::AwaitPending(mut done) => {
let _ = done.changed().await;
}
}
}
}
fn resolve_template_source(
&self,
request: &mut sandbox_v1::CreateSandboxRequest,
) -> Result<templates::TemplateSource, SandboxError> {
match template::Template::parse(&request.template) {
Ok(template::Template::Default) => return Ok(templates::TemplateSource::Default),
Ok(template::Template::DockerImage(image)) => {
return Ok(templates::TemplateSource::DockerImage(image));
}
Err(_) => {}
}
let resolved = self
.manager
.get_template(request.template.trim())
.map_err(SandboxError::from)?;
templates::validate_template_overrides(request)?;
request.template = resolved.canonical_ref();
templates::merge_template_defaults(request, &resolved.entry.defaults);
Ok(templates::TemplateSource::Catalog(Box::new(resolved)))
}
async fn create_once(
&self,
request: sandbox_v1::CreateSandboxRequest,
source: templates::TemplateSource,
create_key: &str,
) -> Result<sandbox_v1::CreateSandboxResponse, SandboxError> {
if !request.id.is_empty()
&& let Some((id, ip_address)) = self
.manager
.replay_sandbox_create(&request.id, create_key)
.await
.map_err(SandboxError::from)?
{
return Ok(sandbox_v1::CreateSandboxResponse {
id,
ip_address,
state: sandbox_v1::SandboxState::Starting.into(),
..Default::default()
});
}
let mut spec = proto_to_spec(request);
if !spec.mounts.is_empty() {
return Err(SandboxError::Unsupported(
"mounts are not supported in Sandbox V1; copy files in with \
WriteFile (`abctl sandbox cp`) instead"
.into(),
));
}
if spec.ssh_public_key.is_some() {
return Err(SandboxError::Unsupported(
"ssh_public_key is not supported in Sandbox V1; use executions \
for interactive access"
.into(),
));
}
spec.rootfs = self.resolve_template(&source).await?;
if let templates::TemplateSource::Catalog(resolved) = &source {
spec.template_warm =
resolved
.entry
.warm
.as_ref()
.map(|warm| arcbox_computer_runtime::TemplateWarmRef {
snapshot_id: warm.snapshot_id.clone(),
vcpus: warm.vcpus,
memory_mib: warm.memory_mib,
});
spec.ready_probe
.clone_from(&resolved.entry.defaults.ready_probe);
}
let (id, ip_address) = self
.manager
.create_sandbox_keyed(spec, create_key)
.await
.map_err(SandboxError::from)?;
register_sandbox_dns(&id, &ip_address);
Ok(sandbox_v1::CreateSandboxResponse {
id,
ip_address,
state: sandbox_v1::SandboxState::Starting.into(),
..Default::default()
})
}
async fn resolve_template(
&self,
source: &templates::TemplateSource,
) -> Result<String, SandboxError> {
match source {
templates::TemplateSource::Default => {
self.rootfs
.ensure_default_rootfs(&self.default_rootfs)
.await
.map_err(|e| SandboxError::Internal(format!("default template: {e}")))?;
Ok(self.default_rootfs.clone())
}
templates::TemplateSource::DockerImage(image) => {
let layout = template::export_docker_image(image)
.await
.map_err(|e| SandboxError::Internal(format!("template {image}: {e:#}")))?;
let pinned = self
.manager
.pinned_rootfs_paths()
.map_err(SandboxError::from)?;
self.rootfs
.convert_layer_to_rootfs(&layout, &pinned)
.await
.map_err(|e| SandboxError::Internal(format!("template {image}: {e}")))
}
templates::TemplateSource::Catalog(resolved) => {
let path = &resolved.entry.rootfs_path;
if !std::path::Path::new(path).is_file() {
return Err(match self.manager.get_template(&resolved.name) {
Err(VmmError::TemplateNotFound(_)) => {
SandboxError::from(VmmError::TemplateNotFound(format!(
"{} (deleted while this create was resolving it)",
resolved.name
)))
}
_ => SandboxError::Internal(format!(
"template {} rootfs {path} is missing; the catalog pin failed",
resolved.name
)),
});
}
Ok(path.clone())
}
}
}
pub async fn stop_request(
&self,
req: sandbox_v1::StopSandboxRequest,
) -> Result<(), SandboxError> {
self.manager
.stop_sandbox(&req.id, req.timeout_seconds)
.await
.map_err(SandboxError::from)?;
deregister_sandbox_dns(&req.id);
Ok(())
}
pub async fn pause_request(
&self,
req: sandbox_v1::PauseSandboxRequest,
) -> Result<(), SandboxError> {
self.manager
.pause_sandbox(&req.id)
.await
.map_err(SandboxError::from)?;
deregister_sandbox_dns(&req.id);
Ok(())
}
pub async fn resume_request(
&self,
req: arcbox_connect::v1::SandboxResumeCommand,
) -> Result<arcbox_connect::v1::SandboxResumeResponse, SandboxError> {
let reason = if req.reason == arcbox_computer_runtime::pause_reason::AUTO_RESUME {
arcbox_computer_runtime::pause_reason::AUTO_RESUME
} else {
arcbox_computer_runtime::pause_reason::RESUME
};
let ip_address = self
.manager
.resume_sandbox(&req.id, reason)
.await
.map_err(SandboxError::from)?;
if !ip_address.is_empty() {
register_sandbox_dns(&req.id, &ip_address);
}
Ok(arcbox_connect::v1::SandboxResumeResponse {
ip_address,
..Default::default()
})
}
pub async fn set_lifecycle_request(
&self,
req: sandbox_v1::SetLifecycleRequest,
) -> Result<(), SandboxError> {
let update = arcbox_computer_runtime::LifecycleUpdate {
ttl_seconds: req.ttl_seconds,
idle_timeout_seconds: req.idle_timeout_seconds,
on_idle: req
.on_idle
.map(|value| idle_action_to_spec(value.as_known().unwrap_or_default())),
};
self.manager
.set_sandbox_lifecycle(&req.id, update)
.await
.map_err(SandboxError::from)
}
pub async fn remove_request(
&self,
req: sandbox_v1::RemoveSandboxRequest,
) -> Result<(), SandboxError> {
self.manager
.remove_sandbox(&req.id, req.force)
.await
.map_err(SandboxError::from)?;
deregister_sandbox_dns(&req.id);
self.clear_stale_completed_create(&req.id);
Ok(())
}
pub fn inspect(&self, payload: &[u8]) -> Result<sandbox_v1::SandboxInfo, SandboxError> {
let req = sandbox_v1::InspectSandboxRequest::decode_from_slice(payload)
.map_err(|e| SandboxError::Decode(e.to_string()))?;
let info = self
.manager
.inspect_sandbox(&req.id)
.map_err(SandboxError::from)?;
Ok(convert::info_to_proto(info))
}
pub fn list(&self, payload: &[u8]) -> Result<sandbox_v1::ListSandboxesResponse, SandboxError> {
let req = sandbox_v1::ListSandboxesRequest::decode_from_slice(payload)
.map_err(|e| SandboxError::Decode(e.to_string()))?;
let state_filter = convert::state_filter(req.state.as_known().unwrap_or_default());
let labels: std::collections::HashMap<String, String> = req.labels.into_iter().collect();
let summaries = self
.manager
.list_sandboxes(state_filter, &labels)
.map_err(SandboxError::from)?;
let (page, next_page_token) =
convert::paginate(summaries, |s| &s.id, req.page_size, &req.page_token);
Ok(sandbox_v1::ListSandboxesResponse {
sandboxes: page.into_iter().map(convert::summary_to_proto).collect(),
next_page_token,
..Default::default()
})
}
pub fn sandbox_network_identity(
&self,
sandbox_id: &str,
) -> Result<arcbox_computer_runtime::SandboxNetworkIdentity, SandboxError> {
self.manager
.sandbox_network_identity(sandbox_id)
.map_err(SandboxError::from)
}
pub async fn wait_startup_cleanup_complete(&self) {
self.manager.wait_startup_cleanup_complete().await;
}
pub async fn pending_cleanup_ticket(
&self,
sandbox_id: &str,
) -> Result<Option<arcbox_connect::v1::SandboxCleanupTicket>, SandboxError> {
Ok(self
.manager
.pending_network_cleanups()
.await
.map_err(SandboxError::from)?
.into_iter()
.find_map(|(id, token)| {
(id == sandbox_id).then_some(arcbox_connect::v1::SandboxCleanupTicket {
id,
token,
startup: false,
..Default::default()
})
}))
}
pub async fn pending_cleanup_tickets(
&self,
) -> Result<Vec<arcbox_connect::v1::SandboxCleanupTicket>, SandboxError> {
let mut tickets = self
.manager
.pending_network_cleanups()
.await
.map_err(SandboxError::from)?
.into_iter()
.map(|(id, token)| arcbox_connect::v1::SandboxCleanupTicket {
id,
token,
startup: false,
..Default::default()
})
.collect::<Vec<_>>();
if let Some(token) = self
.manager
.startup_cleanup_token()
.await
.map_err(SandboxError::from)?
{
tickets.push(arcbox_connect::v1::SandboxCleanupTicket {
id: STARTUP_CLEANUP_ID.into(),
token,
startup: true,
..Default::default()
});
}
Ok(tickets)
}
pub async fn prepare_cleanup(
&self,
ticket: &arcbox_connect::v1::SandboxCleanupTicket,
) -> Result<std::net::Ipv4Addr, SandboxError> {
if ticket.startup {
if ticket.id != STARTUP_CLEANUP_ID {
return Err(SandboxError::Decode(
"invalid sandbox startup cleanup ticket".into(),
));
}
self.manager
.validate_startup_cleanup(&ticket.token)
.await
.map_err(SandboxError::from)?;
return Ok(std::net::Ipv4Addr::UNSPECIFIED);
}
let lease = self
.manager
.validate_network_cleanup(&ticket.id, &ticket.token)
.await
.map_err(SandboxError::from)?;
match lease.ip {
std::net::IpAddr::V4(ip) => Ok(ip),
std::net::IpAddr::V6(ip) => Err(SandboxError::Internal(format!(
"sandbox {} holds {ip}; host cleanup here is IPv4-only",
ticket.id
))),
}
}
pub async fn finalize_cleanup(
&self,
ticket: &arcbox_connect::v1::SandboxCleanupTicket,
) -> Result<(), SandboxError> {
if ticket.startup {
if ticket.id != STARTUP_CLEANUP_ID {
return Err(SandboxError::Decode(
"invalid sandbox startup cleanup ticket".into(),
));
}
return self
.manager
.finalize_startup_cleanup(&ticket.token)
.await
.map_err(SandboxError::from);
}
self.manager
.finalize_network_cleanup(&ticket.id, &ticket.token)
.await
.map_err(SandboxError::from)?;
deregister_sandbox_dns(&ticket.id);
Ok(())
}
pub(crate) fn clear_stale_completed_create(&self, id: &str) {
self.creates.clear_completed_if(id, || {
completed_create_is_stale(
self.manager
.inspect_sandbox(&id.to_owned())
.map(|info| info.state),
)
});
}
}
fn completed_create_is_stale(state: Result<SandboxState, VmmError>) -> bool {
match state {
Ok(SandboxState::Stopped | SandboxState::Failed) | Err(VmmError::NotFound(_)) => true,
Ok(
SandboxState::Starting
| SandboxState::Ready
| SandboxState::Running
| SandboxState::Stopping
| SandboxState::Pausing
| SandboxState::Paused,
)
| Err(_) => false,
}
}
fn register_sandbox_dns(id: &str, ip: &str) {
let Ok(ipv4) = ip.parse::<std::net::Ipv4Addr>() else {
tracing::warn!(id, ip, "invalid sandbox IP for DNS registration");
return;
};
let registry = crate::dns_server::sandbox_registry();
if let Ok(mut map) = registry.write() {
map.insert(id.to_lowercase(), ipv4);
}
}
fn deregister_sandbox_dns(id: &str) {
let registry = crate::dns_server::sandbox_registry();
if let Ok(mut map) = registry.write() {
map.remove(&id.to_lowercase());
}
}
fn proto_to_spec(req: sandbox_v1::CreateSandboxRequest) -> SandboxSpec {
let (vcpus, memory_mib) = (req.limits.vcpus, req.limits.memory_mib);
let mode = match req.network.mode.as_known().unwrap_or_default() {
sandbox_v1::NetworkMode::None => "none",
sandbox_v1::NetworkMode::Enabled | sandbox_v1::NetworkMode::Unspecified => "tap",
};
SandboxSpec {
id: if req.id.is_empty() {
None
} else {
Some(req.id)
},
labels: req.labels.into_iter().collect(),
kernel: String::new(),
rootfs: String::new(),
boot_args: String::new(),
vcpus,
memory_mib,
cmd: req.cmd,
env: req.env.into_iter().collect(),
working_dir: req.working_dir,
user: req.user,
mounts: req
.mounts
.into_iter()
.map(|m| SandboxMountSpec {
source: m.source,
target: m.target,
readonly: m.readonly,
})
.collect(),
network: SandboxNetworkSpec { mode: mode.into() },
ttl_seconds: req.ttl_seconds,
ssh_public_key: req.ssh_public_key,
idle_timeout_seconds: req.idle_timeout_seconds,
on_idle: idle_action_to_spec(req.on_idle.as_known().unwrap_or_default()),
template_warm: None,
ready_probe: None,
}
}
fn idle_action_to_spec(action: sandbox_v1::IdleAction) -> arcbox_computer_runtime::IdleAction {
match action {
sandbox_v1::IdleAction::Pause => arcbox_computer_runtime::IdleAction::Pause,
sandbox_v1::IdleAction::Kill | sandbox_v1::IdleAction::Unspecified => {
arcbox_computer_runtime::IdleAction::Kill
}
}
}
fn create_uses_network(request: &sandbox_v1::CreateSandboxRequest) -> bool {
!matches!(
request.network.mode.as_known().unwrap_or_default(),
sandbox_v1::NetworkMode::None
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn only_networked_create_waits_for_startup_cleanup() {
assert!(create_uses_network(
&sandbox_v1::CreateSandboxRequest::default()
));
assert!(!create_uses_network(&sandbox_v1::CreateSandboxRequest {
network: sandbox_v1::NetworkSpec {
mode: sandbox_v1::NetworkMode::None.into(),
..Default::default()
}
.into(),
..Default::default()
}));
}
#[test]
fn completed_create_is_stale_only_after_terminal_or_removed_state() {
for state in [
SandboxState::Starting,
SandboxState::Ready,
SandboxState::Running,
SandboxState::Stopping,
] {
assert!(!completed_create_is_stale(Ok(state)));
}
for state in [SandboxState::Stopped, SandboxState::Failed] {
assert!(completed_create_is_stale(Ok(state)));
}
assert!(completed_create_is_stale(Err(VmmError::NotFound(
"removed".into()
))));
assert!(!completed_create_is_stale(Err(VmmError::Config(
"inspect failed".into()
))));
}
}