use crate::error::{EngineError, Result};
use crate::persistence::MachinePersistence;
use crate::vm::{HostNetwork, SharedDirConfig, VmConfig, VmId, VmManager};
#[cfg(target_os = "macos")]
use arcbox_constants::ports::AGENT_PORT;
use arcbox_constants::virtiofs::{MOUNT_PRIVATE, MOUNT_USERS, TAG_ARCBOX, TAG_PRIVATE, TAG_USERS};
use chrono::{DateTime, Utc};
use std::collections::HashMap;
use std::net::IpAddr;
use std::path::PathBuf;
use std::sync::{Arc, RwLock};
#[cfg(target_os = "macos")]
use std::sync::{Mutex, PoisonError};
use std::time::Duration;
pub const DEFAULT_MACHINE_NAME: &str = "default";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MachineState {
Created,
Starting,
Running,
Stopping,
Stopped,
}
#[cfg(target_os = "macos")]
mod serial;
#[cfg(test)]
mod tests;
#[derive(Debug, Clone)]
pub struct MachineInfo {
pub name: String,
pub state: MachineState,
pub vm_id: VmId,
pub cid: Option<u32>,
pub cpus: u32,
pub memory_mb: u64,
pub disk_gb: u64,
pub kernel: Option<String>,
pub cmdline: Option<String>,
pub block_devices: Vec<crate::vm::BlockDeviceConfig>,
pub distro: Option<String>,
pub distro_version: Option<String>,
pub disk_path: Option<PathBuf>,
pub ssh_key_path: Option<PathBuf>,
pub ip_address: Option<String>,
pub bridge_ip_address: Option<String>,
pub backend: arcbox_vmm::VmBackend,
pub nested_virt: bool,
pub created_at: DateTime<Utc>,
pub started_at: Option<DateTime<Utc>>,
pub mounts: Vec<MachineMount>,
}
#[derive(Debug, Clone)]
pub struct MachineRootfs {
pub path: PathBuf,
pub format: String,
pub shim: Option<BootShim>,
}
#[derive(Debug, Clone)]
pub struct BootShim {
pub kernel: PathBuf,
pub rootfs: PathBuf,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct MachineMount {
pub host_path: String,
pub guest_path: String,
pub read_only: bool,
}
#[derive(Debug, Clone)]
pub struct MachineConfig {
pub name: String,
pub cpus: u32,
pub memory_mb: u64,
pub disk_gb: u64,
pub kernel: Option<String>,
pub cmdline: Option<String>,
pub block_devices: Vec<crate::vm::BlockDeviceConfig>,
pub rootfs: Option<MachineRootfs>,
pub mounts: Vec<MachineMount>,
pub distro: Option<String>,
pub distro_version: Option<String>,
pub backend: arcbox_vmm::VmBackend,
pub enable_rosetta: bool,
pub nested_virt: bool,
}
impl Default for MachineConfig {
fn default() -> Self {
Self {
name: "default".to_string(),
cpus: arcbox_hypervisor::default_vm_cpu_count(),
memory_mb: 4096,
disk_gb: 50,
kernel: None,
cmdline: None,
block_devices: Vec::new(),
rootfs: None,
mounts: Vec::new(),
distro: None,
distro_version: None,
backend: arcbox_vmm::VmBackend::default(),
enable_rosetta: false,
nested_virt: false,
}
}
}
const fn boot_console() -> &'static str {
#[cfg(target_arch = "x86_64")]
{
"ttyS0"
}
#[cfg(not(target_arch = "x86_64"))]
{
"hvc0"
}
}
const KEEP_KERNEL_NIC_NAMES: &str = "net.ifnames=0";
const QUIET_KERNEL_CONSOLE: &str = "loglevel=4";
fn default_distro_cmdline(rootfs_format: &str) -> String {
let console = boot_console();
format!(
"console={console} root=/dev/vda ro rootfstype={rootfs_format} earlycon \
{KEEP_KERNEL_NIC_NAMES} {QUIET_KERNEL_CONSOLE}"
)
}
fn machine_shim_cmdline(hostname: &str, rootfs_format: &str, mounts: &[MachineMount]) -> String {
use arcbox_constants::cmdline::{
MACHINE_DATA_KEY, MACHINE_INIT_PATH, MACHINE_MOUNTS_KEY, MACHINE_NAME_KEY,
MACHINE_ROOTFS_KEY, MACHINE_ROOTFS_TYPE_KEY,
};
let console = boot_console();
let mut cmdline = format!(
"console={console} root=/dev/vda ro rootfstype=erofs earlycon \
{KEEP_KERNEL_NIC_NAMES} {QUIET_KERNEL_CONSOLE} init={MACHINE_INIT_PATH} \
{MACHINE_ROOTFS_KEY}/dev/vdb {MACHINE_ROOTFS_TYPE_KEY}{rootfs_format} \
{MACHINE_DATA_KEY}/dev/vdc {MACHINE_NAME_KEY}{hostname}"
);
if !mounts.is_empty() {
let table = mounts
.iter()
.enumerate()
.map(|(i, m)| {
let ro = if m.read_only { ":ro" } else { "" };
format!("{}={}{ro}", mount_tag(i), m.guest_path)
})
.collect::<Vec<_>>()
.join(",");
cmdline.push(' ');
cmdline.push_str(MACHINE_MOUNTS_KEY);
cmdline.push_str(&table);
}
cmdline
}
fn mount_tag(index: usize) -> String {
format!("m{index}")
}
pub fn machine_hostname(name: &str) -> Result<String> {
let hostname: String = name
.chars()
.map(|c| if c == '_' || c == '.' { '-' } else { c })
.collect();
let is_label = !hostname.is_empty()
&& hostname.len() <= 63
&& hostname
.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'-')
&& !hostname.starts_with('-')
&& !hostname.ends_with('-');
if is_label {
Ok(hostname)
} else {
Err(EngineError::config(format!(
"machine name '{name}' cannot be a hostname: use 1-63 letters, digits, \
hyphens, underscores and dots, not starting or ending with a separator"
)))
}
}
fn validate_mount(mount: &MachineMount) -> Result<()> {
if !std::path::Path::new(&mount.host_path).is_dir() {
return Err(EngineError::config(format!(
"mount host path '{}' is not a directory",
mount.host_path
)));
}
if !mount.guest_path.starts_with('/') {
return Err(EngineError::config(format!(
"mount guest path '{}' must be absolute",
mount.guest_path
)));
}
if mount.guest_path.contains(',') || mount.guest_path.contains('=') {
return Err(EngineError::config(format!(
"mount guest path '{}' must not contain ',' or '='",
mount.guest_path
)));
}
Ok(())
}
pub struct MachineManager {
machines: RwLock<HashMap<String, MachineInfo>>,
vm_manager: Arc<VmManager>,
persistence: MachinePersistence,
data_dir: PathBuf,
machines_dir: PathBuf,
host_network: HostNetwork,
event_bus: crate::event::EventBus,
#[cfg(target_os = "macos")]
serial_drains: Mutex<HashMap<String, serial::DrainHandle>>,
}
impl MachineManager {
#[must_use]
pub fn new(
vm_manager: Arc<VmManager>,
data_dir: PathBuf,
host_network: HostNetwork,
event_bus: crate::event::EventBus,
) -> Self {
let machines_dir = data_dir.join("machines");
let persistence = MachinePersistence::new(&machines_dir);
let mut shared_dirs = vec![SharedDirConfig::new(
data_dir.to_string_lossy().to_string(),
TAG_ARCBOX,
)];
let users_dir = std::path::Path::new(MOUNT_USERS);
if users_dir.is_dir() {
shared_dirs.push(SharedDirConfig::new(MOUNT_USERS, TAG_USERS));
}
let private_dir = std::path::Path::new(MOUNT_PRIVATE);
if private_dir.is_dir() {
shared_dirs.push(SharedDirConfig::new(MOUNT_PRIVATE, TAG_PRIVATE));
}
let mut machines = HashMap::new();
for persisted in persistence.load_all() {
let needs_recovery = persisted.state.needs_recovery();
if needs_recovery {
tracing::warn!(
"Machine '{}' was running when daemon stopped — marking as stopped",
persisted.name
);
}
let mut vm_shared_dirs = shared_dirs.clone();
for (i, mount) in persisted.mounts.iter().enumerate() {
let mut share = SharedDirConfig::new(mount.host_path.clone(), mount_tag(i));
share.read_only = mount.read_only;
vm_shared_dirs.push(share);
}
let vm_config = VmConfig {
cpus: persisted.cpus,
memory_mb: persisted.memory_mb,
kernel: persisted.kernel.clone(),
cmdline: persisted.cmdline.clone(),
shared_dirs: vm_shared_dirs,
block_devices: persisted.block_devices.clone(),
backend: persisted.backend,
nested_virt: persisted.nested_virt,
..Default::default()
};
if let Ok(vm_id) = vm_manager.create(vm_config) {
let info = MachineInfo {
name: persisted.name.clone(),
state: persisted.state.into(),
vm_id,
cid: None, cpus: persisted.cpus,
memory_mb: persisted.memory_mb,
disk_gb: persisted.disk_gb,
kernel: persisted.kernel.clone(),
cmdline: persisted.cmdline,
block_devices: persisted.block_devices.clone(),
distro: persisted.distro.clone(),
distro_version: persisted.distro_version.clone(),
disk_path: persisted.disk_path.clone().map(PathBuf::from),
ssh_key_path: persisted.ssh_key_path.clone().map(PathBuf::from),
ip_address: persisted.ip_address.clone(),
bridge_ip_address: persisted.bridge_ip_address.clone(),
backend: persisted.backend,
nested_virt: persisted.nested_virt,
created_at: persisted.created_at,
started_at: persisted.started_at,
mounts: persisted.mounts.clone(),
};
machines.insert(persisted.name.clone(), info);
}
if needs_recovery {
if let Err(e) = persistence.update_state(&persisted.name, MachineState::Stopped) {
tracing::warn!(
"Failed to persist corrected state for '{}': {}",
persisted.name,
e
);
}
}
}
tracing::info!("Loaded {} persisted machines", machines.len());
Self {
machines: RwLock::new(machines),
vm_manager,
persistence,
data_dir,
machines_dir,
host_network,
event_bus,
#[cfg(target_os = "macos")]
serial_drains: Mutex::new(HashMap::new()),
}
}
fn publish_event(&self, name: &str, event: crate::event::Event) {
if name != DEFAULT_MACHINE_NAME {
self.event_bus.publish(event);
}
}
pub async fn create(&self, config: MachineConfig) -> Result<String> {
let hostname = machine_hostname(&config.name)?;
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
if machines.contains_key(&config.name) {
return Err(EngineError::already_exists(config.name));
}
if let Some(other) = machines
.keys()
.find(|existing| machine_hostname(existing).ok().as_deref() == Some(hostname.as_str()))
{
return Err(EngineError::already_exists(format!(
"machine name '{}' would take hostname '{hostname}', which machine '{other}' \
already has",
config.name
)));
}
let machine_dir = self.machines_dir.join(&config.name);
std::fs::create_dir_all(&machine_dir)?;
let mut shared_dirs = vec![SharedDirConfig::new(
self.data_dir.to_string_lossy().to_string(),
TAG_ARCBOX,
)];
let users_dir = std::path::Path::new(MOUNT_USERS);
if users_dir.is_dir() {
shared_dirs.push(SharedDirConfig::new(MOUNT_USERS, TAG_USERS));
}
let private_dir = std::path::Path::new(MOUNT_PRIVATE);
if private_dir.is_dir() {
shared_dirs.push(SharedDirConfig::new(MOUNT_PRIVATE, TAG_PRIVATE));
}
if !config.mounts.is_empty() {
if config.rootfs.as_ref().is_none_or(|r| r.shim.is_none()) {
return Err(EngineError::config(
"mounts require a shim-booted distro machine",
));
}
for (i, mount) in config.mounts.iter().enumerate() {
validate_mount(mount)?;
let mut share = SharedDirConfig::new(mount.host_path.clone(), mount_tag(i));
share.read_only = mount.read_only;
shared_dirs.push(share);
}
}
let (kernel, block_devices, cmdline, disk_path) = match &config.rootfs {
Some(rootfs) => {
if rootfs.shim.is_none() && config.kernel.is_none() {
return Err(EngineError::config(
"a distro rootfs without a boot shim requires an explicit \
kernel (pass --kernel)",
));
}
let data_disk = machine_dir.join("data.img");
crate::vm::ensure_sparse_block_image(
&data_disk,
config.disk_gb.saturating_mul(1024 * 1024 * 1024),
)?;
let mut devices = Vec::new();
if let Some(shim) = &rootfs.shim {
devices.push(crate::vm::BlockDeviceConfig {
path: shim.rootfs.to_string_lossy().into_owned(),
read_only: true,
});
}
devices.push(crate::vm::BlockDeviceConfig {
path: rootfs.path.to_string_lossy().into_owned(),
read_only: true,
});
devices.push(crate::vm::BlockDeviceConfig {
path: data_disk.to_string_lossy().into_owned(),
read_only: false,
});
devices.extend(config.block_devices.clone());
let kernel = config.kernel.clone().or_else(|| {
rootfs
.shim
.as_ref()
.map(|s| s.kernel.to_string_lossy().into_owned())
});
let cmdline = config.cmdline.clone().or_else(|| {
Some(match &rootfs.shim {
Some(_) => machine_shim_cmdline(&hostname, &rootfs.format, &config.mounts),
None => default_distro_cmdline(&rootfs.format),
})
});
(kernel, devices, cmdline, Some(data_disk))
}
None => (
config.kernel.clone(),
config.block_devices.clone(),
config.cmdline.clone(),
None,
),
};
let vm_config = VmConfig {
cpus: config.cpus,
memory_mb: config.memory_mb,
kernel: kernel.clone(),
cmdline: cmdline.clone(),
shared_dirs,
block_devices: block_devices.clone(),
rosetta: config.enable_rosetta,
nested_virt: config.nested_virt,
backend: config.backend,
..Default::default()
};
let vm_id = self.vm_manager.create(vm_config)?;
let info = MachineInfo {
name: config.name.clone(),
state: MachineState::Created,
vm_id,
cid: None,
cpus: config.cpus,
memory_mb: config.memory_mb,
disk_gb: config.disk_gb,
kernel,
cmdline,
block_devices,
distro: config.distro,
distro_version: config.distro_version,
disk_path,
ssh_key_path: None,
ip_address: None,
bridge_ip_address: None,
backend: config.backend,
nested_virt: config.nested_virt,
created_at: Utc::now(),
started_at: None,
mounts: config.mounts,
};
self.persistence.save(&info)?;
machines.insert(config.name.clone(), info);
self.publish_event(
&config.name,
crate::event::Event::MachineCreated {
name: config.name.clone(),
},
);
Ok(config.name)
}
pub async fn start(self: &Arc<Self>, name: &str) -> Result<()> {
let (vm_id, cid) = self.assign_cid_for_start(name)?;
let is_machine_vm = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?
.get(name)
.and_then(|m| m.distro.as_ref())
.is_some();
self.vm_manager.start(&vm_id, self.host_network.clone())?;
let started_at = Utc::now();
{
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
if let Some(machine) = machines.get_mut(name) {
machine.state = MachineState::Running;
machine.cid = Some(cid);
machine.ip_address = None;
machine.bridge_ip_address = None;
machine.started_at = Some(started_at);
tracing::info!("Machine '{}' started with CID {}", name, cid);
}
}
#[cfg(target_os = "macos")]
self.start_serial_drain(name, &vm_id);
if let Err(e) = self.persistence.update(name, |m| {
m.state = MachineState::Running.into();
m.ip_address = None;
m.bridge_ip_address = None;
m.started_at = Some(started_at);
}) {
tracing::warn!("Failed to persist state for machine '{}': {}", name, e);
}
if is_machine_vm {
self.wait_for_machine_ready(name).await.map_err(|e| {
EngineError::Machine(format!(
"Machine '{name}' started but readiness check failed: {e}"
))
})?;
}
self.publish_event(
name,
crate::event::Event::MachineStarted {
name: name.to_string(),
},
);
Ok(())
}
async fn wait_for_machine_ready(&self, name: &str) -> Result<()> {
const PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
const INITIAL_DELAY_MS: u64 = 50;
const MAX_DELAY_MS: u64 = 500;
tracing::info!("Waiting for machine '{}' agent to become ready...", name);
let deadline = std::time::Instant::now() + PROBE_TIMEOUT;
let mut delay_ms = INITIAL_DELAY_MS;
let mut attempt: u32 = 0;
let addresses = loop {
attempt += 1;
let probed = match self.connect_agent(name) {
Ok(agent) if agent.is_blocking() => {
tokio::task::block_in_place(|| probe_ip_blocking(agent, name, attempt))?
}
Ok(agent) => probe_ip_async(agent, name, attempt).await?,
Err(e) => {
tracing::trace!("Machine '{}' connect failed (attempt {attempt}): {e}", name);
None
}
};
if let Some(addresses) = probed {
break addresses;
}
if std::time::Instant::now() >= deadline {
self.log_console_tail(name);
return Err(EngineError::Machine(format!(
"Machine '{name}' agent did not report a routable IP within timeout"
)));
}
tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await;
delay_ms = (delay_ms * 3 / 2).min(MAX_DELAY_MS);
};
{
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
if let Some(machine) = machines.get_mut(name) {
machine.ip_address = Some(addresses.ip.clone());
machine.bridge_ip_address.clone_from(&addresses.bridge_ip);
}
}
if let Err(e) = self.persistence.update(name, |m| {
m.ip_address = Some(addresses.ip.clone());
m.bridge_ip_address.clone_from(&addresses.bridge_ip);
}) {
tracing::warn!("Failed to persist IP for machine '{}': {}", name, e);
}
tracing::info!(
machine = name,
ip = %addresses.ip,
bridge_ip = addresses.bridge_ip.as_deref().unwrap_or("none"),
"machine ready"
);
Ok(())
}
fn assign_cid_for_start(&self, name: &str) -> Result<(VmId, u32)> {
let (vm_id, cid) = {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state == MachineState::Running {
return Err(EngineError::invalid_state(format!(
"machine '{name}' is already running"
)));
}
if machine.state == MachineState::Starting || machine.state == MachineState::Stopping {
return Err(EngineError::invalid_state(format!(
"machine '{name}' is in transition state"
)));
}
let used: std::collections::HashSet<u32> =
machines.values().filter_map(|m| m.cid).collect();
let cid = (3..=u32::MAX)
.find(|c| !used.contains(c))
.expect("fewer than u32::MAX machines");
(machine.vm_id.clone(), cid)
};
self.vm_manager.set_guest_cid(&vm_id, cid)?;
Ok((vm_id, cid))
}
#[must_use]
pub fn vm_manager(&self) -> &VmManager {
&self.vm_manager
}
#[cfg(all(target_os = "macos", feature = "vmnet"))]
pub fn vmnet_interface_mac(&self, name: &str) -> Option<String> {
let machines = self.machines.read().ok()?;
let machine = machines.get(name)?;
self.vm_manager.vmnet_interface_mac(&machine.vm_id)
}
pub fn bridge_mac(&self, name: &str) -> Option<String> {
let machines = self.machines.read().ok()?;
let machine = machines.get(name)?;
Some(crate::vm::bridge_nic_mac_for_vm_id(&machine.vm_id))
}
#[must_use]
pub fn get_cid(&self, name: &str) -> Option<u32> {
self.machines.read().ok()?.get(name)?.cid
}
#[cfg(target_os = "macos")]
pub fn connect_agent(&self, name: &str) -> Result<crate::agent_client::AgentClient> {
use crate::agent_client::AgentClient;
let (cid, vm_id) = {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state != MachineState::Running {
return Err(EngineError::invalid_state(format!(
"machine '{name}' is not running"
)));
}
let cid = machine
.cid
.ok_or_else(|| EngineError::invalid_state("CID not assigned"))?;
(cid, machine.vm_id.clone())
};
let backend = self.vm_manager.backend(&vm_id)?;
let fd = self.connect_vsock_port(name, AGENT_PORT)?;
match backend {
arcbox_vmm::VmBackend::Hv => AgentClient::from_fd_blocking(cid, fd),
arcbox_vmm::VmBackend::Vz => AgentClient::from_fd_async(cid, fd),
}
}
#[cfg(target_os = "macos")]
pub fn connect_vsock_port(&self, name: &str, port: u32) -> Result<std::os::unix::io::RawFd> {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state != MachineState::Running {
return Err(EngineError::invalid_state(format!(
"machine '{name}' is not running"
)));
}
self.vm_manager.connect_vsock(&machine.vm_id, port)
}
#[cfg(target_os = "linux")]
pub fn connect_agent(&self, name: &str) -> Result<crate::agent_client::AgentClient> {
use crate::agent_client::AgentClient;
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state != MachineState::Running {
return Err(EngineError::invalid_state(format!(
"machine '{}' is not running",
name
)));
}
let cid = machine
.cid
.ok_or_else(|| EngineError::invalid_state("CID not assigned"))?;
Ok(AgentClient::new(cid))
}
#[cfg(target_os = "linux")]
pub fn connect_vsock_port(&self, name: &str, port: u32) -> Result<std::os::unix::io::RawFd> {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state != MachineState::Running {
return Err(EngineError::invalid_state(format!(
"machine '{}' is not running",
name
)));
}
self.vm_manager.connect_vsock(&machine.vm_id, port)
}
pub async fn ping_agent(self: Arc<Self>, machine_name: String) -> Result<()> {
let manager = Arc::clone(&self);
let name = machine_name.clone();
let connected = tokio::task::spawn_blocking(move || manager.connect_agent(&name)).await;
match connected {
Ok(Ok(mut agent)) => {
if agent.is_blocking() {
tokio::task::spawn_blocking(move || agent.ping_blocking().map(|_| ()))
.await
.unwrap_or_else(|e| {
Err(EngineError::Vm(format!("agent ping task panicked: {e}")))
})
} else {
agent.ping().await.map(|_| ())
}
}
Ok(Err(e)) => Err(e),
Err(e) => Err(EngineError::Vm(format!("agent connect task panicked: {e}"))),
}
}
pub async fn record_bridge_address(
self: Arc<Self>,
machine_name: String,
) -> Result<Option<String>> {
let manager = Arc::clone(&self);
let name = machine_name.clone();
let connected = tokio::task::spawn_blocking(move || manager.connect_agent(&name)).await;
let info = match connected {
Ok(Ok(mut agent)) => {
if agent.is_blocking() {
tokio::task::spawn_blocking(move || agent.get_system_info_blocking())
.await
.unwrap_or_else(|e| {
Err(EngineError::Vm(format!("system info task panicked: {e}")))
})?
} else {
agent.get_system_info().await?
}
}
Ok(Err(e)) => return Err(e),
Err(e) => {
return Err(EngineError::Vm(format!("agent connect task panicked: {e}")));
}
};
let bridge_ip = Some(info.bridge_ip_address).filter(|ip| !ip.is_empty());
{
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
if let Some(machine) = machines.get_mut(&machine_name) {
machine.bridge_ip_address.clone_from(&bridge_ip);
}
}
if let Err(e) = self.persistence.update(&machine_name, |m| {
m.bridge_ip_address.clone_from(&bridge_ip);
}) {
tracing::warn!(
"Failed to persist bridge address for machine '{}': {}",
machine_name,
e
);
}
Ok(bridge_ip)
}
pub async fn trim_disk(self: Arc<Self>, machine_name: String) -> Result<u64> {
let manager = Arc::clone(&self);
let name = machine_name.clone();
let connected = tokio::task::spawn_blocking(move || manager.connect_agent(&name)).await;
let response = match connected {
Ok(Ok(mut agent)) => {
if agent.is_blocking() {
tokio::task::spawn_blocking(move || agent.disk_trim_blocking())
.await
.unwrap_or_else(|e| {
Err(EngineError::Vm(format!("disk trim task panicked: {e}")))
})?
} else {
agent.disk_trim().await?
}
}
Ok(Err(e)) => return Err(e),
Err(e) => return Err(EngineError::Vm(format!("agent connect task panicked: {e}"))),
};
tracing::debug!(machine = %machine_name, result = %response.result, "disk trim done");
Ok(response.bytes_trimmed)
}
#[cfg(target_os = "macos")]
fn start_serial_drain(&self, name: &str, vm_id: &VmId) {
let handle = match self.vm_manager.dup_serial_readers(vm_id) {
Ok(Some(readers)) => serial::spawn(name, readers),
Ok(None) => {
tracing::debug!(machine = name, "no host console pipes to drain");
return;
}
Err(e) => {
tracing::warn!(
machine = name,
"console pipes unavailable, not draining: {e}"
);
return;
}
};
self.serial_drains
.lock()
.unwrap_or_else(PoisonError::into_inner)
.insert(name.to_owned(), handle);
}
#[cfg(target_os = "macos")]
fn stop_serial_drain(&self, name: &str) {
self.serial_drains
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remove(name);
}
#[cfg(not(target_os = "macos"))]
fn stop_serial_drain(&self, _name: &str) {}
#[cfg(target_os = "macos")]
fn log_console_tail(&self, name: &str) {
let kept = self
.serial_drains
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(name)
.map(serial::DrainHandle::tail);
let (console, agent_log) = match kept {
Some(tail) => tail,
None => {
let vm_id = match self.machines.read() {
Ok(machines) => machines.get(name).map(|m| m.vm_id.clone()),
Err(_) => None,
};
let Some(vm_id) = vm_id else { return };
let last_lines = |output: Result<String>| -> Vec<String> {
let Ok(text) = output else { return Vec::new() };
let mut lines: Vec<String> =
text.lines().rev().take(40).map(str::to_owned).collect();
lines.reverse();
lines
};
(
last_lines(self.vm_manager.read_console_output(&vm_id)),
last_lines(self.vm_manager.read_agent_log_output(&vm_id)),
)
}
};
for (label, lines) in [("console", console), ("agent-log", agent_log)] {
for line in lines {
tracing::warn!(machine = %name, "machine {label}: {line}");
}
}
}
#[cfg(not(target_os = "macos"))]
fn log_console_tail(&self, _name: &str) {}
pub fn debug_snapshot(&self, name: &str) -> Result<arcbox_vmm::VmDebugSnapshot> {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
self.vm_manager.debug_snapshot(&machine.vm_id)
}
pub fn stop(&self, name: &str) -> Result<()> {
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get_mut(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if !matches!(
machine.state,
MachineState::Running | MachineState::Stopping
) {
return Err(EngineError::invalid_state(format!(
"machine '{name}' is not running"
)));
}
#[cfg(target_os = "macos")]
self.vm_manager
.force_stop_without_hypervisor(&machine.vm_id)?;
#[cfg(not(target_os = "macos"))]
self.vm_manager.stop(&machine.vm_id)?;
machine.state = MachineState::Stopped;
machine.cid = None;
machine.bridge_ip_address = None;
self.stop_serial_drain(name);
if let Err(e) = self.persist_stopped(name) {
tracing::warn!(
"Failed to persist stopped state for machine '{}': {}",
name,
e
);
}
self.publish_event(
name,
crate::event::Event::MachineStopped {
name: name.to_string(),
},
);
Ok(())
}
#[must_use]
pub fn vm_self_stopped(&self, name: &str) -> Option<bool> {
let machines = self.machines.read().ok()?;
let vm_id = machines.get(name)?.vm_id.clone();
drop(machines);
self.vm_manager.vm_self_stopped(&vm_id)
}
pub fn reboot(&self, name: &str) -> Result<()> {
let vm_id = {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?
.vm_id
.clone()
};
self.vm_manager.reboot(&vm_id)?;
tracing::info!("Rebooted machine '{}'", name);
Ok(())
}
pub fn set_backend(&self, name: &str, backend: arcbox_vmm::VmBackend) -> Result<()> {
let vm_id = {
let machines = self
.machines
.read()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if matches!(
machine.state,
MachineState::Running | MachineState::Starting
) {
return Err(EngineError::invalid_state(format!(
"cannot switch backend while machine '{name}' is {:?}",
machine.state
)));
}
machine.vm_id.clone()
};
self.vm_manager.set_backend(&vm_id, backend)?;
if let Some(machine) = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?
.get_mut(name)
{
machine.backend = backend;
}
self.persistence.update(name, |m| m.backend = backend)?;
Ok(())
}
pub fn graceful_stop(&self, name: &str, timeout: Duration) -> Result<bool> {
let vm_id = {
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get_mut(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state != MachineState::Running {
return Err(EngineError::invalid_state(format!(
"machine '{name}' is not running"
)));
}
machine.state = MachineState::Stopping;
machine.vm_id.clone()
};
match self.vm_manager.graceful_stop(&vm_id, timeout) {
Ok(true) => {
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get_mut(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
machine.state = MachineState::Stopped;
machine.cid = None;
machine.bridge_ip_address = None;
self.stop_serial_drain(name);
if let Err(e) = self.persist_stopped(name) {
tracing::warn!(
"Failed to persist stopped state for machine '{}': {}",
name,
e
);
}
drop(machines);
self.publish_event(
name,
crate::event::Event::MachineStopped {
name: name.to_string(),
},
);
Ok(true)
}
Ok(false) => {
self.rollback_stopping(name);
Ok(false)
}
Err(e) => {
self.rollback_stopping(name);
Err(e)
}
}
}
fn persist_stopped(&self, name: &str) -> Result<()> {
self.persistence.update(name, |m| {
m.state = MachineState::Stopped.into();
m.bridge_ip_address = None;
})
}
fn rollback_stopping(&self, name: &str) {
if let Ok(mut machines) = self.machines.write() {
if let Some(machine) = machines.get_mut(name) {
if machine.state == MachineState::Stopping {
machine.state = MachineState::Running;
}
}
}
}
#[must_use]
pub fn get(&self, name: &str) -> Option<MachineInfo> {
self.machines.read().ok()?.get(name).cloned()
}
#[must_use]
pub fn exists(&self, name: &str) -> bool {
self.machines
.read()
.is_ok_and(|machines| machines.contains_key(name))
}
#[must_use]
pub fn list(&self) -> Vec<MachineInfo> {
self.machines
.read()
.map(|m| m.values().cloned().collect())
.unwrap_or_default()
}
pub fn remove(&self, name: &str, force: bool) -> Result<()> {
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
let machine = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
if machine.state == MachineState::Running && !force {
return Err(EngineError::invalid_state(
"cannot remove running machine (use --force)".to_string(),
));
}
if machine.state == MachineState::Running {
let vm_id = machine.vm_id.clone();
drop(machines); self.vm_manager.stop(&vm_id)?;
self.stop_serial_drain(name);
machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
}
let vm_id = {
let m = machines
.get(name)
.ok_or_else(|| EngineError::not_found(name.to_string()))?;
m.vm_id.clone()
};
self.vm_manager.remove(&vm_id)?;
machines.remove(name);
self.persistence.remove(name)?;
drop(machines);
self.publish_event(
name,
crate::event::Event::MachineRemoved {
name: name.to_string(),
},
);
tracing::info!("Removed machine '{}'", name);
Ok(())
}
#[cfg(target_os = "macos")]
pub fn take_inbound_listener_manager(
&self,
name: &str,
) -> Option<arcbox_net::darwin::inbound_relay::InboundListenerManager> {
let vm_id = {
let machines = self.machines.read().ok()?;
let machine = machines.get(name)?;
if machine.state != MachineState::Running {
return None;
}
machine.vm_id.clone()
};
self.vm_manager.take_inbound_listener_manager(&vm_id)
}
pub fn register_mock_machine(&self, name: &str, cid: u32) -> Result<()> {
let mut machines = self
.machines
.write()
.map_err(|_| EngineError::LockPoisoned)?;
if machines.contains_key(name) {
return Ok(()); }
let info = MachineInfo {
name: name.to_string(),
state: MachineState::Running,
vm_id: VmId::new(), cid: Some(cid),
cpus: arcbox_hypervisor::default_vm_cpu_count(),
memory_mb: 4096,
disk_gb: 50,
kernel: None,
cmdline: None,
block_devices: Vec::new(),
distro: None,
distro_version: None,
disk_path: None,
ssh_key_path: None,
ip_address: None,
bridge_ip_address: None,
backend: arcbox_vmm::VmBackend::default(),
nested_virt: false,
created_at: Utc::now(),
started_at: None,
mounts: Vec::new(),
};
machines.insert(name.to_string(), info);
tracing::debug!("Registered mock machine '{}' with CID {}", name, cid);
Ok(())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct GuestAddresses {
ip: String,
bridge_ip: Option<String>,
}
fn readiness_addresses(
info: &arcbox_connect::v1::SystemInfo,
name: &str,
attempt: u32,
) -> Option<GuestAddresses> {
if info.distro_init_pending {
tracing::trace!(
"Machine '{name}' distro init still starting (attempt {attempt}); not ready",
);
return None;
}
let ip = select_routable_ip(&info.ip_addresses)?;
let bridge_ip = Some(info.bridge_ip_address.clone()).filter(|ip| !ip.is_empty());
Some(GuestAddresses { ip, bridge_ip })
}
fn probe_ip_blocking(
mut agent: crate::agent_client::AgentClient,
name: &str,
attempt: u32,
) -> Result<Option<GuestAddresses>> {
let resp = match agent.ping_blocking() {
Ok(resp) => resp,
Err(e) => {
tracing::trace!("Machine '{}' ping failed (attempt {attempt}): {e}", name);
return Ok(None);
}
};
crate::agent_client::AgentClient::check_agent_protocol(&resp)?;
tracing::debug!(
"Machine '{}' agent reachable (version: {}, attempt {})",
name,
resp.version,
attempt,
);
match agent.get_system_info_blocking() {
Ok(info) => Ok(readiness_addresses(&info, name, attempt)),
Err(e) => {
tracing::trace!(
"Machine '{}' get_system_info failed (attempt {attempt}): {e}",
name,
);
Ok(None)
}
}
}
async fn probe_ip_async(
mut agent: crate::agent_client::AgentClient,
name: &str,
attempt: u32,
) -> Result<Option<GuestAddresses>> {
let resp = match agent.ping().await {
Ok(resp) => resp,
Err(e) => {
tracing::trace!("Machine '{}' ping failed (attempt {attempt}): {e}", name);
return Ok(None);
}
};
crate::agent_client::AgentClient::check_agent_protocol(&resp)?;
tracing::debug!(
"Machine '{}' agent reachable (version: {}, attempt {})",
name,
resp.version,
attempt,
);
match agent.get_system_info().await {
Ok(info) => Ok(readiness_addresses(&info, name, attempt)),
Err(e) => {
tracing::trace!(
"Machine '{}' get_system_info failed (attempt {attempt}): {e}",
name,
);
Ok(None)
}
}
}
fn select_routable_ip(ips: &[String]) -> Option<String> {
let mut ipv6_candidate = None;
for ip in ips {
let Ok(addr) = ip.parse::<IpAddr>() else {
continue;
};
if addr.is_loopback() || addr.is_multicast() || addr.is_unspecified() {
continue;
}
match addr {
IpAddr::V4(v4) => return Some(v4.to_string()),
IpAddr::V6(v6) => {
if v6.is_unicast_link_local() {
continue;
}
if ipv6_candidate.is_none() {
ipv6_candidate = Some(v6.to_string());
}
}
}
}
ipv6_candidate
}