mod control;
mod filesystem;
mod icon;
mod kubernetes;
mod machine;
#[cfg(target_os = "macos")]
mod macos;
mod migration;
mod process;
mod sandbox_errors;
mod sandbox_resume;
mod snapshot;
mod stats;
mod stream_input;
mod system;
mod template;
use std::sync::Arc;
use std::time::Duration;
use std::sync::OnceLock;
use arcbox_core::Runtime;
use arcbox_core::vm_lifecycle::DEFAULT_MACHINE_NAME;
use connectrpc::{ConnectError, RequestContext};
use tokio_stream::{Stream, StreamExt as _};
pub use control::SandboxServiceImpl;
pub use filesystem::SandboxFilesystemServiceImpl;
pub use icon::IconServiceImpl;
pub use kubernetes::KubernetesServiceImpl;
pub use machine::MachineServiceImpl;
#[cfg(target_os = "macos")]
pub use macos::MacosServiceImpl;
pub use migration::MigrationServiceImpl;
pub use process::SandboxProcessServiceImpl;
pub use snapshot::SandboxSnapshotServiceImpl;
pub use stats::StatsServiceImpl;
pub use system::{SetupState, SystemServiceImpl};
pub use template::TemplateServiceImpl;
pub type SharedRuntime = Arc<OnceLock<Arc<Runtime>>>;
const KEEPALIVE_INTERVAL: Duration = Duration::from_secs(15);
fn with_keepalive<S, T>(
stream: S,
keepalive: fn() -> T,
) -> impl Stream<Item = Result<T, ConnectError>>
where
S: Stream<Item = Result<T, ConnectError>>,
{
stream
.timeout(KEEPALIVE_INTERVAL)
.map(move |item| item.unwrap_or_else(|_elapsed| Ok(keepalive())))
}
#[cfg(target_os = "macos")]
pub(crate) async fn run_macos_blocking<T, Fut, F>(f: F) -> Result<T, ConnectError>
where
T: Send + 'static,
Fut: std::future::Future<Output = arcbox_core::Result<T>>,
F: FnOnce() -> Fut + Send + 'static,
{
tokio::task::spawn_blocking(move || {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| ConnectError::internal(format!("macOS runtime: {e}")))?
.block_on(f())
.map_err(|e| ConnectError::internal(e.to_string()))
})
.await
.map_err(|e| ConnectError::internal(format!("macOS task join: {e}")))?
}
pub(crate) trait ConnectRuntimeExt {
fn ready(&self) -> Result<&Arc<Runtime>, ConnectError>;
fn ready_for_write(&self, machine: &str) -> Result<&Arc<Runtime>, ConnectError>;
}
impl ConnectRuntimeExt for SharedRuntime {
fn ready(&self) -> Result<&Arc<Runtime>, ConnectError> {
self.get()
.ok_or_else(|| ConnectError::unavailable("daemon is starting, runtime not ready yet"))
}
fn ready_for_write(&self, machine: &str) -> Result<&Arc<Runtime>, ConnectError> {
let runtime = self.ready()?;
runtime
.ensure_storage_writes_available(machine)
.map_err(crate::ApiError::from)?;
Ok(runtime)
}
}
pub(crate) trait ContextExt {
fn machine_id(&self) -> Result<String, ConnectError>;
fn sandbox_machine_id(&self) -> Result<String, ConnectError>;
}
impl ContextExt for RequestContext {
fn machine_id(&self) -> Result<String, ConnectError> {
match self.header("x-machine") {
None => Ok(DEFAULT_MACHINE_NAME.to_owned()),
Some(value) => match value.to_str() {
Ok("") => Ok(DEFAULT_MACHINE_NAME.to_owned()),
Ok(s) => Ok(s.to_string()),
Err(_) => Err(ConnectError::invalid_argument(
"invalid x-machine header: must be valid UTF-8",
)),
},
}
}
fn sandbox_machine_id(&self) -> Result<String, ConnectError> {
let machine = self.machine_id()?;
if machine != DEFAULT_MACHINE_NAME {
return Err(ConnectError::invalid_argument(
"Sandbox V1 is available only on the System VM",
));
}
Ok(machine)
}
}
fn port_protocol(
protocol: arcbox_connect::sandbox_v1::PortProtocol,
) -> arcbox_computer::ports::SandboxPortProtocol {
use arcbox_computer::ports::SandboxPortProtocol;
match protocol {
arcbox_connect::sandbox_v1::PortProtocol::Udp => SandboxPortProtocol::Udp,
_ => SandboxPortProtocol::Tcp,
}
}
fn exposed_port(
mapping: arcbox_core::SandboxPortMapping,
) -> arcbox_connect::sandbox_v1::ExposedPort {
use arcbox_connect::sandbox_v1::{ExposedPort, PortProtocol};
use arcbox_core::SandboxPortProtocol;
ExposedPort {
sandbox_port: u32::from(mapping.sandbox_port),
host_port: u32::from(mapping.host_port),
protocol: match mapping.protocol {
SandboxPortProtocol::Tcp => PortProtocol::Tcp,
SandboxPortProtocol::Udp => PortProtocol::Udp,
}
.into(),
..Default::default()
}
}
#[must_use]
pub fn router(runtime: SharedRuntime) -> connectrpc::Router {
let clone = || Arc::clone(&runtime);
let sandbox_operations = Arc::new(arcbox_computer::locks::SandboxOperationLocks::default());
let router = connectrpc::Router::new()
.add_service(Arc::new(SandboxServiceImpl::new(
clone(),
Arc::clone(&sandbox_operations),
)))
.add_service(Arc::new(SandboxProcessServiceImpl::new(
clone(),
Arc::clone(&sandbox_operations),
)))
.add_service(Arc::new(SandboxFilesystemServiceImpl::new(
clone(),
Arc::clone(&sandbox_operations),
)))
.add_service(Arc::new(SandboxSnapshotServiceImpl::new(
clone(),
Arc::clone(&sandbox_operations),
)))
.add_service(Arc::new(TemplateServiceImpl::new(clone())))
.add_service(Arc::new(IconServiceImpl::new()))
.add_service(Arc::new(StatsServiceImpl::new(clone())))
.add_service(Arc::new(KubernetesServiceImpl::new(clone())))
.add_service(Arc::new(MigrationServiceImpl::new(clone())))
.add_service(Arc::new(MachineServiceImpl::new(clone())));
#[cfg(target_os = "macos")]
let router = router.add_service(Arc::new(MacosServiceImpl::new(clone())));
router
}
#[must_use]
pub fn router_with_system(runtime: SharedRuntime, system: SystemServiceImpl) -> connectrpc::Router {
router(runtime).add_service(Arc::new(system))
}
#[cfg(test)]
mod tests {
use super::*;
pub(super) fn storage_runtime() -> (std::path::PathBuf, SharedRuntime) {
let directory = std::env::temp_dir().join(format!(
"arcbox-api-storage-admission-{}",
uuid::Uuid::new_v4()
));
std::fs::create_dir_all(&directory).unwrap();
let runtime = Arc::new(
Runtime::new(arcbox_core::config::Config {
data_dir: directory.clone(),
..Default::default()
})
.unwrap(),
);
(directory, Arc::new(OnceLock::from(runtime)))
}
#[tokio::test]
async fn storage_admission_rejects_reserved_and_held_writes_but_preserves_reads() {
let (directory, shared) = storage_runtime();
let runtime = shared.ready().unwrap();
for machine in [DEFAULT_MACHINE_NAME, "rosetta", "dev"] {
assert!(shared.ready_for_write(machine).is_ok());
}
let reservation = runtime.machine_manager().reserve_storage().unwrap();
for machine in [DEFAULT_MACHINE_NAME, "rosetta"] {
let error = shared.ready_for_write(machine).err().unwrap();
assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
}
assert!(shared.ready().is_ok());
assert!(shared.ready_for_write("dev").is_ok());
drop(reservation);
assert!(shared.ready_for_write(DEFAULT_MACHINE_NAME).is_ok());
let hold = runtime.machine_manager().storage_hold_path();
std::fs::create_dir_all(hold.parent().unwrap()).unwrap();
std::fs::write(&hold, "offline-check").unwrap();
assert!(!runtime.storage_writes_protected());
for machine in [DEFAULT_MACHINE_NAME, "rosetta"] {
let error = shared.ready_for_write(machine).err().unwrap();
assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
}
assert!(shared.ready().is_ok());
assert!(shared.ready_for_write("dev").is_ok());
std::fs::remove_file(hold).unwrap();
assert!(shared.ready_for_write(DEFAULT_MACHINE_NAME).is_ok());
std::fs::remove_dir_all(directory).unwrap();
}
#[tokio::test]
async fn failed_hold_write_keeps_public_storage_writes_protected() {
let (directory, shared) = storage_runtime();
let runtime = shared.ready().unwrap();
let hold = runtime.machine_manager().storage_hold_path();
std::fs::create_dir_all(&hold).unwrap();
assert!(
runtime
.recover_storage(arcbox_connect::v1::recover_storage_request::Action::CheckOnly)
.await
.is_err()
);
std::fs::remove_dir(&hold).unwrap();
assert!(!runtime.storage_recovery_active());
assert!(runtime.storage_writes_protected());
assert!(
runtime
.machine_manager()
.ensure_storage_available(DEFAULT_MACHINE_NAME)
.is_err()
);
for machine in [DEFAULT_MACHINE_NAME, "rosetta"] {
let error = shared.ready_for_write(machine).err().unwrap();
assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
}
assert!(shared.ready().is_ok());
assert!(shared.ready_for_write("dev").is_ok());
assert!(matches!(
runtime.trim_machine_disk(DEFAULT_MACHINE_NAME).await,
Err(arcbox_core::CoreError::Common(
arcbox_error::CommonError::InvalidState(_)
))
));
std::fs::remove_dir_all(directory).unwrap();
}
#[tokio::test]
async fn resume_rechecks_storage_admission_after_waiting_for_the_operation_lock() {
let (directory, shared) = storage_runtime();
let runtime = shared.ready().unwrap();
let operations = arcbox_computer::locks::SandboxOperationLocks::default();
let operation = operations.lock(DEFAULT_MACHINE_NAME, "sandbox").await;
let resumed = sandbox_resume::resume(
runtime,
&operations,
DEFAULT_MACHINE_NAME,
"sandbox",
sandbox_resume::REASON_RESUME,
);
tokio::pin!(resumed);
tokio::select! {
biased;
result = &mut resumed => panic!("resume must wait for its operation lock: {result:?}"),
() = tokio::task::yield_now() => {}
}
let reservation = runtime.machine_manager().reserve_storage().unwrap();
drop(operation);
let error = resumed.await.unwrap_err();
assert_eq!(error.code, connectrpc::ErrorCode::FailedPrecondition);
drop(reservation);
std::fs::remove_dir_all(directory).unwrap();
}
fn ctx_with(header: Option<&str>) -> RequestContext {
let mut headers = http::HeaderMap::new();
if let Some(value) = header {
headers.insert("x-machine", value.parse().expect("valid header value"));
}
RequestContext::new(headers)
}
#[test]
fn machine_id_defaults_to_the_system_vm() {
assert_eq!(
ctx_with(None).machine_id().expect("absent header is valid"),
DEFAULT_MACHINE_NAME
);
assert_eq!(
ctx_with(Some(""))
.machine_id()
.expect("empty header is valid"),
DEFAULT_MACHINE_NAME
);
assert_eq!(
ctx_with(Some("other-vm"))
.machine_id()
.expect("explicit header is valid"),
"other-vm"
);
}
#[test]
fn sandbox_machine_id_accepts_only_the_system_vm() {
assert_eq!(DEFAULT_MACHINE_NAME, "default");
for header in [None, Some(""), Some(DEFAULT_MACHINE_NAME)] {
assert_eq!(
ctx_with(header)
.sandbox_machine_id()
.expect("System VM routing should be accepted"),
DEFAULT_MACHINE_NAME
);
}
let error = ctx_with(Some("other-vm"))
.sandbox_machine_id()
.expect_err("other machines must be rejected");
assert_eq!(error.code, connectrpc::ErrorCode::InvalidArgument);
}
#[test]
fn exposed_port_preserves_the_authoritative_mapping() {
let port = exposed_port(arcbox_core::SandboxPortMapping {
sandbox_port: 8080,
host_port: 45_000,
protocol: arcbox_core::SandboxPortProtocol::Udp,
});
assert_eq!(port.sandbox_port, 8080);
assert_eq!(port.host_port, 45_000);
assert_eq!(
port.protocol.as_known(),
Some(arcbox_connect::sandbox_v1::PortProtocol::Udp)
);
}
}