From d17c83af6c64652b7512c0e701f3d60f04d92fce Mon Sep 17 00:00:00 2001 From: Kevin Wang Date: Mon, 17 Aug 2026 01:40:38 -0700 Subject: [PATCH 1/2] refactor(vmm): name bridge preparation explicitly --- dstack/vmm/src/app.rs | 5 +++-- dstack/vmm/src/netd.rs | 37 +++++++++++++++++++++++++--------- dstack/vmm/src/vm_launcher.rs | 38 ++--------------------------------- 3 files changed, 32 insertions(+), 48 deletions(-) diff --git a/dstack/vmm/src/app.rs b/dstack/vmm/src/app.rs index da9e6c445..ef7a6589c 100644 --- a/dstack/vmm/src/app.rs +++ b/dstack/vmm/src/app.rs @@ -6,7 +6,8 @@ use crate::{ config::{Config, NetworkFilterMode, Networking, NetworkingMode, ProcessAnnotation, Protocol}, logrotate, netd::{ - self, InterfaceIdentity, PrepareMacvtapRequest, PrepareRequest, Request as NetdRequest, + self, InterfaceIdentity, PrepareBridgeRequest, PrepareMacvtapRequest, + Request as NetdRequest, }, }; @@ -563,7 +564,7 @@ impl App { nic_index, ); let request = match network.mode { - NetworkingMode::Bridge => NetdRequest::Prepare(PrepareRequest { + NetworkingMode::Bridge => NetdRequest::PrepareBridge(PrepareBridgeRequest { identity: identity.clone(), bridge: network.bridge.clone(), mac, diff --git a/dstack/vmm/src/netd.rs b/dstack/vmm/src/netd.rs index 79ffdde3a..458a0638f 100644 --- a/dstack/vmm/src/netd.rs +++ b/dstack/vmm/src/netd.rs @@ -49,7 +49,7 @@ pub struct InterfaceIdentity { } #[derive(Debug, Clone, Serialize, Deserialize)] -pub struct PrepareRequest { +pub struct PrepareBridgeRequest { #[serde(flatten)] pub identity: InterfaceIdentity, pub bridge: String, @@ -73,7 +73,7 @@ pub struct PrepareMacvtapRequest { #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(tag = "operation", rename_all = "snake_case")] pub enum Request { - Prepare(PrepareRequest), + PrepareBridge(PrepareBridgeRequest), PrepareMacvtap(PrepareMacvtapRequest), Remove { #[serde(flatten)] @@ -121,7 +121,7 @@ pub struct PreparedInterface { pub async fn request(socket: &Path, request: &Request) -> Result { let operation = match request { - Request::Prepare(_) => "prepare", + Request::PrepareBridge(_) => "prepare_bridge", Request::PrepareMacvtap(_) => "prepare_macvtap", Request::Remove { .. } => "remove", Request::Check { .. } => "check", @@ -248,8 +248,8 @@ async fn read_request(stream: &mut UnixStream) -> Result { fn handle_request(libvirt_uri: &str, request: Request) -> Result<(String, Option)> { let _lock = OperationLock::acquire()?; match request { - Request::Prepare(request) => { - prepare_interface(libvirt_uri, &request).map(|tap| (tap, None)) + Request::PrepareBridge(request) => { + prepare_bridge(libvirt_uri, &request).map(|tap| (tap, None)) } Request::PrepareMacvtap(request) => prepare_macvtap( libvirt_uri, @@ -362,8 +362,8 @@ impl Drop for OperationLock { } } -fn prepare_interface(libvirt_uri: &str, request: &PrepareRequest) -> Result { - validate_prepare(request)?; +fn prepare_bridge(libvirt_uri: &str, request: &PrepareBridgeRequest) -> Result { + validate_prepare_bridge(request)?; let tap = tap_name(&request.identity); // A failed VMM start may leave a deterministic resource behind. Replacing // it makes prepare idempotent without accepting a caller-selected TAP. @@ -430,7 +430,7 @@ fn delete_binding(uri: &str, tap: &str) -> Result<()> { bail!("virsh failed to delete binding {tap}: {}", error.trim()) } -fn binding_xml(request: &PrepareRequest, tap: &str) -> String { +fn binding_xml(request: &PrepareBridgeRequest, tap: &str) -> String { let owner_uuid = stable_uuid(&request.identity); let owner_name = format!( "dstack:{}:{}:{}", @@ -472,7 +472,7 @@ fn stable_uuid(identity: &InterfaceIdentity) -> Uuid { Uuid::from_bytes(bytes) } -fn validate_prepare(request: &PrepareRequest) -> Result<()> { +fn validate_prepare_bridge(request: &PrepareBridgeRequest) -> Result<()> { validate_identity(&request.identity)?; validate_name("bridge", &request.bridge, 15, "_.-")?; if !Path::new("/sys/class/net") @@ -670,7 +670,7 @@ mod tests { #[test] fn binding_xml_escapes_values() { - let request = PrepareRequest { + let request = PrepareBridgeRequest { identity: identity("instance<&", "vm", 0), bridge: "br0".into(), mac: "02:00:00:00:00:01".into(), @@ -719,6 +719,23 @@ mod tests { assert!(value.get("identity").is_none()); } + #[test] + fn bridge_prepare_has_a_dedicated_operation() { + let request = Request::PrepareBridge(PrepareBridgeRequest { + identity: identity("instance", "vm", 0), + bridge: "br0".into(), + mac: "02:00:00:00:00:01".into(), + qemu_uid: 1000, + filter: "clean-traffic".into(), + parameters: BTreeMap::new(), + }); + let value = serde_json::to_value(request).unwrap(); + assert_eq!(value["operation"], "prepare_bridge"); + assert_eq!(value["instance_id"], "instance"); + assert_eq!(value["bridge"], "br0"); + assert!(value.get("identity").is_none()); + } + #[tokio::test] async fn disconnected_client_is_confined_to_one_connection() { let (mut server, client) = UnixStream::pair().unwrap(); diff --git a/dstack/vmm/src/vm_launcher.rs b/dstack/vmm/src/vm_launcher.rs index 4ce1bebba..eaff26440 100644 --- a/dstack/vmm/src/vm_launcher.rs +++ b/dstack/vmm/src/vm_launcher.rs @@ -71,12 +71,6 @@ impl Drop for SocketCleanup { } } -impl LaunchSpec { - fn is_single_process(&self) -> bool { - self.swtpm.is_none() - } -} - struct PreparedOpenFiles { // Retain the original device handles through the fork or exec handoff; // they close automatically when this prepared set leaves scope. @@ -120,7 +114,7 @@ fn prepare_open_files(open_files: &[OpenFile]) -> Result { }) } -fn prepare_exec(expected_parent: libc::pid_t, mappings: &[(i32, i32)]) -> std::io::Result<()> { +fn prepare_child(expected_parent: libc::pid_t, mappings: &[(i32, i32)]) -> std::io::Result<()> { setpgid(Pid::from_raw(0), Pid::from_raw(0))?; prctl::set_pdeathsig(Signal::SIGKILL)?; if getppid().as_raw() != expected_parent { @@ -153,23 +147,13 @@ fn spawn_child(spec: &ChildCommand, open_files: &[OpenFile]) -> Result { unsafe { command .as_std_mut() - .pre_exec(move || prepare_exec(parent, &mappings)); + .pre_exec(move || prepare_child(parent, &mappings)); } command .spawn() .with_context(|| format!("failed to start {}", spec.command)) } -fn exec_in_place(spec: &ChildCommand, open_files: &[OpenFile]) -> Result<()> { - let prepared = prepare_open_files(open_files)?; - let parent = getppid().as_raw(); - prepare_exec(parent, &prepared.mappings).context("failed to prepare in-place QEMU exec")?; - let mut command = std::process::Command::new(&spec.command); - command.args(&spec.args); - let error = command.exec(); - Err(error).with_context(|| format!("failed to exec {}", spec.command)) -} - async fn stop_child(child: &mut Child, name: &str, grace: Duration) { let Some(pid) = child.id() else { return; @@ -217,9 +201,6 @@ pub async fn run(spec_path: &Path) -> Result<()> { let raw = fs_err::read(spec_path) .with_context(|| format!("failed to read launch spec {}", spec_path.display()))?; let spec: LaunchSpec = serde_json::from_slice(&raw).context("failed to parse launch spec")?; - if spec.is_single_process() { - return exec_in_place(&spec.qemu, &spec.open_files); - } let _socket_cleanup = spec.swtpm_socket.clone().map(SocketCleanup); if let Some(socket) = &spec.swtpm_socket { if socket.exists() { @@ -369,21 +350,6 @@ mod tests { Ok(()) } - #[test] - fn single_process_requires_no_swtpm() { - let mut spec = LaunchSpec { - qemu: shell("exit 0".into()), - swtpm: None, - swtpm_socket: None, - open_files: vec![], - startup_timeout_ms: 1_000, - shutdown_timeout_ms: 1_000, - }; - assert!(spec.is_single_process()); - spec.swtpm = Some(shell("exit 0".into())); - assert!(!spec.is_single_process()); - } - #[tokio::test] async fn qemu_failure_stops_and_reaps_swtpm() { let dir = tempfile::tempdir().unwrap(); From b9ba18afb3d0cd9a9c4ce0a176c2b747d41e747d Mon Sep 17 00:00:00 2001 From: Kevin Wang Date: Mon, 17 Aug 2026 01:41:26 -0700 Subject: [PATCH 2/2] feat(vmm): exec single-process launches in place --- dstack/vmm/src/vm_launcher.rs | 38 +++++++++++++++++++++++++++++++++-- 1 file changed, 36 insertions(+), 2 deletions(-) diff --git a/dstack/vmm/src/vm_launcher.rs b/dstack/vmm/src/vm_launcher.rs index eaff26440..4ce1bebba 100644 --- a/dstack/vmm/src/vm_launcher.rs +++ b/dstack/vmm/src/vm_launcher.rs @@ -71,6 +71,12 @@ impl Drop for SocketCleanup { } } +impl LaunchSpec { + fn is_single_process(&self) -> bool { + self.swtpm.is_none() + } +} + struct PreparedOpenFiles { // Retain the original device handles through the fork or exec handoff; // they close automatically when this prepared set leaves scope. @@ -114,7 +120,7 @@ fn prepare_open_files(open_files: &[OpenFile]) -> Result { }) } -fn prepare_child(expected_parent: libc::pid_t, mappings: &[(i32, i32)]) -> std::io::Result<()> { +fn prepare_exec(expected_parent: libc::pid_t, mappings: &[(i32, i32)]) -> std::io::Result<()> { setpgid(Pid::from_raw(0), Pid::from_raw(0))?; prctl::set_pdeathsig(Signal::SIGKILL)?; if getppid().as_raw() != expected_parent { @@ -147,13 +153,23 @@ fn spawn_child(spec: &ChildCommand, open_files: &[OpenFile]) -> Result { unsafe { command .as_std_mut() - .pre_exec(move || prepare_child(parent, &mappings)); + .pre_exec(move || prepare_exec(parent, &mappings)); } command .spawn() .with_context(|| format!("failed to start {}", spec.command)) } +fn exec_in_place(spec: &ChildCommand, open_files: &[OpenFile]) -> Result<()> { + let prepared = prepare_open_files(open_files)?; + let parent = getppid().as_raw(); + prepare_exec(parent, &prepared.mappings).context("failed to prepare in-place QEMU exec")?; + let mut command = std::process::Command::new(&spec.command); + command.args(&spec.args); + let error = command.exec(); + Err(error).with_context(|| format!("failed to exec {}", spec.command)) +} + async fn stop_child(child: &mut Child, name: &str, grace: Duration) { let Some(pid) = child.id() else { return; @@ -201,6 +217,9 @@ pub async fn run(spec_path: &Path) -> Result<()> { let raw = fs_err::read(spec_path) .with_context(|| format!("failed to read launch spec {}", spec_path.display()))?; let spec: LaunchSpec = serde_json::from_slice(&raw).context("failed to parse launch spec")?; + if spec.is_single_process() { + return exec_in_place(&spec.qemu, &spec.open_files); + } let _socket_cleanup = spec.swtpm_socket.clone().map(SocketCleanup); if let Some(socket) = &spec.swtpm_socket { if socket.exists() { @@ -350,6 +369,21 @@ mod tests { Ok(()) } + #[test] + fn single_process_requires_no_swtpm() { + let mut spec = LaunchSpec { + qemu: shell("exit 0".into()), + swtpm: None, + swtpm_socket: None, + open_files: vec![], + startup_timeout_ms: 1_000, + shutdown_timeout_ms: 1_000, + }; + assert!(spec.is_single_process()); + spec.swtpm = Some(shell("exit 0".into())); + assert!(!spec.is_single_process()); + } + #[tokio::test] async fn qemu_failure_stops_and_reaps_swtpm() { let dir = tempfile::tempdir().unwrap();