use std::fmt::Write;
use std::{
path::{Path, PathBuf},
process::Stdio,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
time::Duration,
};
use base64ct::{Base64, Encoding};
use dir_lock::DirLock;
use eyre::{Context, OptionExt, Result, bail};
use serde::{Deserialize, Serialize};
use tokio::fs;
use tokio::io::{AsyncBufReadExt, AsyncRead};
use tokio::process::{Child, Command};
use tracing::{Instrument, debug, debug_span, error, info, instrument, trace};
use crate::types::PmemMount;
use crate::vm_images::VmImage;
use crate::vms::QEMU_PID_FILENAME;
use crate::{
runner::CancellationTokens,
types::{BindMount, PublishPort},
utils::ExecutablePaths,
};
pub(crate) fn command_as_string(cmd: &Command) -> String {
let program_str = cmd.as_std().get_program().to_string_lossy();
let args_str = cmd
.as_std()
.get_args()
.map(|x| x.to_string_lossy())
.map(|x| {
if x.contains(' ') {
format!("\"{x}\"")
} else {
format!("{x}")
}
})
.collect::<Vec<_>>()
.join(" ");
format!("{program_str} {args_str}")
}
pub(crate) async fn extract_kernel(virt_copy_out_path: &Path, vm_image: &VmImage) -> Result<()> {
let dest_dir = vm_image.image_path.parent().ok_or_eyre(format!(
"Image {:?} doesn't have a parent",
vm_image.image_path
))?;
let mut virt_copy_out_cmd = Command::new(virt_copy_out_path);
let files_to_extract = if let Some(initrd) = &vm_image.initrd_path {
vec![
format!("/boot/{}", initrd.file_name().unwrap().to_string_lossy()),
format!(
"/boot/{}",
vm_image.kernel_path.file_name().unwrap().to_string_lossy()
),
]
} else {
vec![format!(
"/boot/{}",
vm_image.kernel_path.file_name().unwrap().to_string_lossy()
)]
};
virt_copy_out_cmd
.args(["-a", &vm_image.image_path.to_string_lossy()])
.args(files_to_extract)
.arg(dest_dir);
let virt_copy_out_cmd_str = command_as_string(&virt_copy_out_cmd);
info!("Extracting kernel from {:?}", vm_image.image_path);
debug!("{virt_copy_out_cmd_str}");
let virt_copy_out_output = virt_copy_out_cmd.output().await?;
if !virt_copy_out_output.status.success() {
bail!(
"virt_copy_out failed: {}",
String::from_utf8_lossy(&virt_copy_out_output.stderr)
);
}
Ok(())
}
#[instrument]
pub(crate) async fn convert_ovmf_uefi_variables(
vm_dir: &Path,
source_image: &Path,
) -> Result<PathBuf> {
let output_file = vm_dir.join("OVMF_VARS.4m.fd.qcow2");
let mut qemu_img_cmd = Command::new("qemu-img");
qemu_img_cmd
.arg("convert")
.args(["-O", "qcow2"])
.arg(source_image)
.arg(&output_file);
let qemu_img_cmd_str = command_as_string(&qemu_img_cmd);
info!("Converting OVMF UEFI vars file to qcow2");
debug!("{qemu_img_cmd_str}");
let qemu_img_output = qemu_img_cmd.output().await?;
if !qemu_img_output.status.success() {
bail!(
"qemu-img convert failed: {}",
String::from_utf8_lossy(&qemu_img_output.stderr)
);
}
Ok(output_file)
}
#[instrument]
pub(crate) async fn create_overlay_image(source_image: &Path, overlay_image: &Path) -> Result<()> {
let source_image_str = source_image.to_string_lossy();
let backing_file = format!("backing_file={source_image_str},backing_fmt=qcow2,nocow=on");
let mut qemu_img_cmd = Command::new("qemu-img");
qemu_img_cmd
.arg("create")
.args(["-o", &backing_file])
.args(["-f", "qcow2"])
.arg(overlay_image);
let qemu_img_cmd_str = command_as_string(&qemu_img_cmd);
info!("Creating overlay image");
debug!("{qemu_img_cmd_str}");
let qemu_img_output = qemu_img_cmd.output().await?;
if !qemu_img_output.status.success() {
bail!(
"qemu-img create failed: {}",
String::from_utf8_lossy(&qemu_img_output.stderr)
);
}
Ok(())
}
fn log_child_output(stream: impl AsyncRead + Unpin + Send + 'static) {
let span = debug_span!(parent: None, "virtofsd process");
tokio::spawn(
async move {
let reader = tokio::io::BufReader::new(stream);
let mut lines = reader.lines();
while let Ok(Some(line)) = lines.next_line().await {
debug!("{line}");
}
}
.instrument(span),
);
}
#[instrument]
pub(crate) async fn launch_virtiofsd(
virtiofsd_path: &Path,
vm_dir: &Path,
volume: &BindMount,
) -> Result<Child> {
let socket_path = vm_dir.join(volume.socket_name());
let mut virtiofsd_cmd = Command::new("unshare");
virtiofsd_cmd
.arg("--map-root-user")
.arg("--map-auto")
.arg("--")
.arg(virtiofsd_path)
.args(["--shared-dir", &volume.source.to_string_lossy()])
.args(["--socket-path", &socket_path.to_string_lossy()])
.args(["--cache", "never"])
.arg("--allow-direct-io")
.arg("--allow-mmap")
.args(["--thread-pool-size", "8"])
.args(["--sandbox", "chroot"]);
if volume.read_only {
virtiofsd_cmd.arg("--readonly");
}
let virtiofsd_cmd_str = command_as_string(&virtiofsd_cmd);
info!("Running virtiofsd for share '{volume}'");
trace!("{virtiofsd_cmd_str}");
let mut virtiofsd_child = virtiofsd_cmd
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true)
.spawn()?;
tokio::select! {
_ = tokio::time::sleep(Duration::from_millis(250)) => {},
_ = virtiofsd_child.wait() => {
error!("virtiofsd process exited early, that's usually a bad sign");
let virtiofsd_output = virtiofsd_child.wait_with_output().await?;
bail!("virtiofsd failed: {}", String::from_utf8(virtiofsd_output.stderr)?);
}
}
if let Some(stdout) = virtiofsd_child.stdout.take() {
log_child_output(stdout);
}
if let Some(stderr) = virtiofsd_child.stderr.take() {
log_child_output(stderr);
}
Ok(virtiofsd_child)
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct QemuLaunchOpts {
pub volumes: Vec<BindMount>,
pub pmems: Vec<PmemMount>,
pub published_ports: Vec<PublishPort>,
pub vm_image: VmImage,
pub ovmf_uefi_vars_path: PathBuf,
pub show_vm_window: bool,
pub pubkey: String,
pub cid: u32,
pub is_warmup: bool,
pub disable_kvm: bool,
pub memory: u64,
}
#[instrument(skip(cancellation_tokens, lock, tool_paths, qemu_launch_opts))]
pub(crate) async fn launch_qemu(
cancellation_tokens: CancellationTokens,
qemu_should_exit: Arc<AtomicBool>,
vm_dir: &Path,
lock: Option<DirLock>,
tool_paths: ExecutablePaths,
qemu_launch_opts: QemuLaunchOpts,
) -> Result<()> {
let overlay_image_str = qemu_launch_opts.vm_image.image_path.to_string_lossy();
let kernel_path_str = qemu_launch_opts.vm_image.kernel_path.to_string_lossy();
let ovmf_uefi_vars_str = qemu_launch_opts.ovmf_uefi_vars_path.to_string_lossy();
let sysinfo_system = sysinfo::System::new_with_specifics(
sysinfo::RefreshKind::nothing().with_cpu(sysinfo::CpuRefreshKind::everything()),
);
let memory = qemu_launch_opts.memory;
let logical_core_count = sysinfo_system.cpus().len();
let ssh_pubkey_base64 = Base64::encode_string(qemu_launch_opts.pubkey.as_bytes());
let cid = qemu_launch_opts.cid;
let hostfwd: String =
qemu_launch_opts
.published_ports
.iter()
.fold(String::new(), |mut output, p| {
let _ = write!(
output,
",hostfwd=:{}:{}-:{}",
p.host_ip, p.host_port, p.vm_port
);
output
});
let qmp_socket_path = vm_dir.join("qmp.sock,server,wait=off");
let qmp_socket_path_str = qmp_socket_path.to_string_lossy();
let total_pmem_size = qemu_launch_opts.pmems.iter().fold(0, |acc, x| acc + x.size);
let qemu_maxmem = format!(",maxmem={}G", memory + total_pmem_size);
let mut qemu_cmd = Command::new(tool_paths.qemu_path.clone());
qemu_cmd
.args(["-machine", "hpet=off"])
.args(["-smp", &logical_core_count.to_string()])
.args(["-kernel", &kernel_path_str])
.args(["-append", "rw root=/dev/vda3"])
.args(["-device", &format!("vhost-vsock-pci,id=vhost-vsock-pci0,guest-cid={cid}")])
.args(["-nic", &format!("user,model=virtio{hostfwd}")])
.args(["-device", "virtio-balloon,free-page-reporting=on"])
.args(["-m", &format!("{memory}G{qemu_maxmem}")])
.args(["-object", &format!("memory-backend-memfd,id=mem0,merge=on,share=on,size={memory}G")])
.args(["-numa", "node,memdev=mem0"])
.args([
"-drive",
"if=pflash,format=raw,unit=0,file=/usr/share/edk2/x64/OVMF_CODE.4m.fd,readonly=on",
])
.args([
"-drive",
&format!("if=pflash,unit=1,file={ovmf_uefi_vars_str}"),
])
.args(["-drive", &format!("if=virtio,node-name=overlay-disk,file={overlay_image_str}")])
.args(["-qmp", &format!("unix:{qmp_socket_path_str}")])
.args([
"-smbios",
&format!(
"type=11,value=io.systemd.credential.binary:ssh.authorized_keys.root={ssh_pubkey_base64}"
),
]);
if !qemu_launch_opts.disable_kvm {
qemu_cmd.args(["-accel", "kvm"]).args(["-cpu", "host"]);
}
if let Some(ref initrd_path) = qemu_launch_opts.vm_image.initrd_path {
let initrd_path_str = initrd_path.to_string_lossy();
qemu_cmd.args(["-initrd", &initrd_path_str]);
}
let mut virtiofsd_handles = vec![];
if !qemu_launch_opts.is_warmup {
qemu_cmd.arg("-snapshot");
let mut fstab_entries = vec![];
for (i, vol) in qemu_launch_opts.volumes.iter().enumerate() {
let virtiofsd_child = launch_virtiofsd(&tool_paths.virtiofsd_path, vm_dir, vol)
.await
.wrap_err(format!("Failed to launch virtiofsd for {vol}"))?;
virtiofsd_handles.push(virtiofsd_child);
let socket_path = vm_dir.join(vol.socket_name());
let socket_path_str = socket_path.to_string_lossy();
let tag = vol.tag();
let dest_path = vol.dest.to_string_lossy();
let read_only = if vol.read_only {
String::from(",ro")
} else {
String::new()
};
let fstab_entry = format!("{tag} {dest_path} virtiofs defaults{read_only} 0 0");
fstab_entries.push(fstab_entry);
qemu_cmd
.args([
"-chardev",
&format!("socket,id=char{i},path={socket_path_str}"),
])
.args([
"-device",
&format!("vhost-user-fs-pci,chardev=char{i},tag={tag}"),
]);
}
for (i, pmem) in qemu_launch_opts.pmems.iter().enumerate() {
let fstab_entry = format!(
"/dev/pmem{i} {} ext4 rw,relatime,dax=always,x-systemd.makefs 0 0",
pmem.dest.to_string_lossy()
);
fstab_entries.push(fstab_entry);
let pmem_file = vm_dir.join(format!("pmem{i}.pmem"));
qemu_cmd.args([
"-object",
&format!("memory-backend-file,id=pmem{i},share=on,merge=on,discard-data=on,mem-path={},size={}G", pmem_file.to_string_lossy(), pmem.size)]);
qemu_cmd.args([
"-device",
&format!("virtio-pmem-pci,memdev=pmem{i},id=nv{i}"),
]);
}
if !fstab_entries.is_empty() {
let fstab = fstab_entries.join("\n");
let fstab_base64 = Base64::encode_string(fstab.as_bytes());
qemu_cmd.args([
"-smbios",
&format!("type=11,value=io.systemd.credential.binary:fstab.extra={fstab_base64}"),
]);
}
}
if !qemu_launch_opts.show_vm_window {
qemu_cmd.arg("-nographic");
}
let qemu_cmd_str = command_as_string(&qemu_cmd);
info!("Running QEMU");
trace!("{qemu_cmd_str}");
let qemu_child = qemu_cmd
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.kill_on_drop(true)
.spawn()?;
let qemu_pid_path = vm_dir.join(QEMU_PID_FILENAME);
let qemu_pid = qemu_child
.id()
.ok_or_eyre("QEMU has no pid, maybe it exited early?")?
.to_string();
fs::write(&qemu_pid_path, &qemu_pid).await?;
trace!("Writing QEMU pid {qemu_pid} to {qemu_pid_path:?}");
if let Some(lock) = lock {
trace!("Unlocking {:?}", lock.path());
lock.drop_async().await.expect("Couldn't drop lock");
}
let qemu_output = tokio::select! {
_ = cancellation_tokens.qemu.cancelled() => {
debug!("QEMU task was cancelled");
return Ok(());
}
val = qemu_child.wait_with_output() => {
if qemu_should_exit.load(Ordering::SeqCst) {
info!("QEMU has finished running");
return Ok(());
}
error!("QEMU process exited early, that's usually a bad sign");
val?
}
};
if !qemu_output.status.success() {
error!("QEMU failed: {}", String::from_utf8(qemu_output.stderr)?);
cancellation_tokens.ssh.cancel();
}
Ok(())
}