vmm: migration: handle in dedicated thread (make async)

This puts the send-migration action into a dedicated thread, laying the
groundwork for many follow-ups towards first-class live-migration in
CH.

This means:

1. The send-migration call will exit sooner (just trigger the
   migration - dispatch semantics)
2. Other API calls can be triggered while a migration is ongoing but
   will not be able to alter the VM as the VM's ownership is transferred
   from the VMM to the migration thread. Example: hotplugging won't work
   (which is good).
3. This is the basis for migration statistics via a dedicated endpoint
   (future work).

The whole change was done with a special focus on graceful recover and
cleanup: even if anything on the migration paths go wrong, the proper
cleanups are already executed and the VMM can take back the ownership
of the VM.

The receive-migration API call remains blocking. To observe any status
changes about the migration on the sender side, one can observe the
event-monitor output and look for `vm.migration-{failed,finished}`.

These changes are inspired by [0] but differ significantly in details.

[0] https://github.com/cloud-hypervisor/cloud-hypervisor/pull/7038

On-behalf-of: SAP philipp.schuster@sap.com
Signed-off-by: Philipp Schuster <philipp.schuster@cyberus-technology.de>
This commit is contained in:
Philipp Schuster
2026-06-01 20:36:26 +02:00
committed by Bo Chen
parent 5d835bdff4
commit 796fc055bd
2 changed files with 218 additions and 88 deletions

View File

@@ -63,6 +63,7 @@ use crate::migration::{recv_vm_config, recv_vm_state};
use crate::migration_transport::{
ReceiveAdditionalConnections, ReceiveListener, SendAdditionalConnections, SocketStream,
};
use crate::migration_worker::{MigrationWorker, MigrationWorkerHandle, MigrationWorkerResult};
use crate::seccomp_filters::{Thread, get_seccomp_filter};
use crate::vm::{Error as VmError, Vm, VmState};
use crate::vm_config::{
@@ -89,7 +90,6 @@ pub mod landlock;
pub mod memory_manager;
pub mod migration;
pub mod migration_transport;
#[expect(unused)]
mod migration_worker;
mod pci_segment;
pub mod seccomp_filters;
@@ -629,21 +629,18 @@ pub struct VmmThreadHandle {
/// Models the current ownership and associated state of the VM from the
/// perspective of the VMM.
#[allow(clippy::large_enum_variant)]
pub enum VmOwnership {
Owned(Vm),
/// The VM is temporarily owned by an ongoing migration worker.
Migration {
migration_worker_handle: MigrationWorkerHandle,
/// Snapshot returned while the VMM cannot inspect the worker-owned VM.
vm_info_response: VmInfoResponse,
},
None,
}
impl VmOwnership {
/// Returns a shared reference to the underlying VM, if available.
fn as_ref(&self) -> Option<&Vm> {
match self {
VmOwnership::Owned(vm) => Some(vm),
_ => None,
}
}
/// Returns a mutable reference to the underlying VM, if available.
fn as_mut(&mut self) -> Option<&mut Vm> {
match self {
@@ -653,13 +650,22 @@ impl VmOwnership {
}
/// Takes the inner VM if it is currently owned.
fn take_owned(&mut self) -> Option<Vm> {
fn take_owned_or(&mut self, none_error: VmError) -> result::Result<Vm, VmError> {
match mem::replace(self, VmOwnership::None) {
VmOwnership::Owned(vm) => Some(vm),
old => {
VmOwnership::Owned(vm) => Ok(vm),
old @ VmOwnership::Migration { .. } => {
*self = old;
None
Err(VmError::VmMigrating)
}
VmOwnership::None => Err(none_error),
}
}
/// Returns an error if the VM is currently migrated.
fn ok_or_migrating(&self) -> result::Result<(), VmError> {
match self {
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
_ => Ok(()),
}
}
}
@@ -1456,7 +1462,10 @@ impl Vmm {
Ok(())
}
/// Performs a migration including all its phases.
/// Performs a migration.
///
/// Runs after-migration cleanup only on success. Callers must handle failed
/// migrations.
fn send_migration(
vm: &mut Vm,
#[cfg(all(feature = "kvm", target_arch = "x86_64"))]
@@ -1703,6 +1712,8 @@ impl Vmm {
prefault: bool,
memory_restore_mode: MemoryRestoreMode,
) -> std::result::Result<(), VmError> {
self.vm.ok_or_migrating()?;
let snapshot = recv_vm_state(source_url).map_err(VmError::Restore)?;
#[cfg(all(feature = "kvm", target_arch = "x86_64"))]
let vm_snapshot = get_vm_snapshot(&snapshot).map_err(VmError::Restore)?;
@@ -1770,8 +1781,64 @@ impl Vmm {
self.vm.as_mut().unwrap().restore()
}
/// Handles the outcome of the migration thread.
fn check_migration(&mut self) {}
/// Handles the outcome of the migration worker thread.
fn check_migration(&mut self) {
let VmOwnership::Migration {
migration_worker_handle,
..
} = mem::replace(&mut self.vm, VmOwnership::None)
else {
panic!("Should only be called after a migration was started");
};
let MigrationWorkerResult {
vm,
migration_result: migration_res,
initial_vm_state,
} = migration_worker_handle.join();
let mut try_resume_vm_after_failed_migration = |mut vm: Vm| {
// A late failure may leave the VM paused.
if initial_vm_state == VmState::Running && vm.get_state() == VmState::Paused {
match vm.resume() {
Ok(_) => {
info!("Resumed VM successfully after failed migration");
}
Err(e) => {
error!("Failed resuming VM after failed migration: {e}");
self.exit_evt.write(1).unwrap();
}
}
}
// Ensure full VM performance. The operation is idempotent.
let _ = vm.stop_dirty_log().inspect_err(|e| {
warn!("Failed stopping dirty log after resuming VM: {e} - VM performance might be slower than usual");
});
self.vm = VmOwnership::Owned(vm);
};
match migration_res {
Ok(()) => {
self.vm = VmOwnership::None;
let mut vm = vm;
// Since the VMM explicitly no longer owns the VM, the exit
// event won't call the shutdown path automatically.
if let Err(e) = vm.shutdown() {
error!("Failed shutting down the VM after migration: {e}");
}
if let Err(e) = self.exit_evt.write(1) {
error!("Failed exiting the VMM after migration: {e}");
}
}
Err(e) => {
error!("Migration failed: {e}");
try_resume_vm_after_failed_migration(vm);
}
}
}
fn control_loop(
&mut self,
@@ -1812,6 +1879,7 @@ impl Vmm {
info!("VM exit event");
// Consume the event.
self.exit_evt.read().map_err(Error::EventFdRead)?;
// TODO: Future follow-up must resolve lifecycle handling while migrating.
self.vmm_shutdown().map_err(Error::VmmShutdown)?;
break 'outer;
@@ -1820,11 +1888,13 @@ impl Vmm {
info!("VM reset event");
// Consume the event.
self.reset_evt.read().map_err(Error::EventFdRead)?;
// TODO: Future follow-up must resolve lifecycle handling while migrating.
self.vm_reboot().map_err(Error::VmReboot)?;
}
EpollDispatch::GuestExit => {
info!("VM guest exit event");
self.guest_exit_evt.read().map_err(Error::EventFdRead)?;
// TODO: Future follow-up must resolve lifecycle handling while migrating.
if self.no_shutdown {
self.vm_shutdown().map_err(Error::VmShutdown)?;
} else {
@@ -1833,8 +1903,9 @@ impl Vmm {
}
}
EpollDispatch::ActivateVirtioDevices => {
// TODO: Future follow-up must resolve virtio activation handling while migrating.
let count = self.activate_evt.read().map_err(Error::EventFdRead)?;
if let VmOwnership::Owned(ref vm) = self.vm {
let count = self.activate_evt.read().map_err(Error::EventFdRead)?;
info!("Trying to activate pending virtio devices: count = {count}");
vm.activate_virtio_devices()
.map_err(Error::ActivateVirtioDevices)?;
@@ -1858,10 +1929,12 @@ impl Vmm {
// Read from the API receiver channel
let gdb_request = gdb_receiver.recv().map_err(Error::GdbRequestRecv)?;
let response = if let VmOwnership::Owned(ref mut vm) = self.vm {
vm.debug_request(&gdb_request.payload, gdb_request.cpu_id)
} else {
Err(VmError::VmNotRunning)
let response = match self.vm {
VmOwnership::Owned(ref mut vm) => {
vm.debug_request(&gdb_request.payload, gdb_request.cpu_id)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
.map_err(gdb::Error::Vm);
@@ -1906,6 +1979,8 @@ fn apply_landlock(vm_config: &mut VmConfig) -> result::Result<(), LandlockError>
impl RequestHandler for Vmm {
fn vm_create(&mut self, config: Box<VmConfig>) -> result::Result<(), VmError> {
self.vm.ok_or_migrating()?;
// We only store the passed VM config.
// The VM will be created when being asked to boot it.
if self.vm_config.is_some() {
@@ -1928,6 +2003,8 @@ impl RequestHandler for Vmm {
}
fn vm_boot(&mut self) -> result::Result<(), VmError> {
self.vm.ok_or_migrating()?;
tracer::start();
info!("Booting VM");
event!("vm", "booting");
@@ -2003,6 +2080,7 @@ impl RequestHandler for Vmm {
fn vm_pause(&mut self) -> result::Result<(), VmError> {
match self.vm {
VmOwnership::Owned(ref mut vm) => vm.pause().map_err(VmError::Pause),
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
@@ -2010,6 +2088,7 @@ impl RequestHandler for Vmm {
fn vm_resume(&mut self) -> result::Result<(), VmError> {
match self.vm {
VmOwnership::Owned(ref mut vm) => vm.resume().map_err(VmError::Resume),
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
@@ -2026,11 +2105,14 @@ impl RequestHandler for Vmm {
.map_err(VmError::SnapshotSend)
})
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
fn vm_restore(&mut self, restore_cfg: RestoreConfig) -> result::Result<(), VmError> {
self.vm.ok_or_migrating()?;
if self.vm_config.is_some() || matches!(self.vm, VmOwnership::Owned(_)) {
return Err(VmError::VmAlreadyCreated);
}
@@ -2095,12 +2177,13 @@ impl RequestHandler for Vmm {
VmOwnership::Owned(ref mut vm) => {
vm.coredump(destination_url).map_err(VmError::Coredump)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
fn vm_shutdown(&mut self) -> result::Result<(), VmError> {
let mut vm = self.vm.take_owned().ok_or(VmError::VmNotRunning)?;
let mut vm = self.vm.take_owned_or(VmError::VmNotRunning)?;
// Drain console_info so that the FDs are not reused
let _ = self.console_info.take();
let r = vm.shutdown();
@@ -2118,7 +2201,7 @@ impl RequestHandler for Vmm {
// Drop VM early to release disk locks and free other resources before
// we reboot.
let config = {
let mut vm = self.vm.take_owned().ok_or(VmError::VmNotCreated)?;
let mut vm = self.vm.take_owned_or(VmError::VmNotCreated)?;
let config = vm.get_config();
// First we stop the current VM
vm.shutdown()?;
@@ -2186,35 +2269,42 @@ impl RequestHandler for Vmm {
}
fn vm_info(&self) -> result::Result<VmInfoResponse, VmError> {
match &self.vm_config {
Some(vm_config) => {
let state = match &self.vm {
VmOwnership::Owned(vm) => vm.get_state(),
VmOwnership::None => VmState::Created,
};
let config = vm_config.lock().unwrap().clone();
let mut memory_actual_size =
config.memory.total_size() - config.memory.hotplugged_size();
if let VmOwnership::Owned(vm) = &self.vm {
memory_actual_size = memory_actual_size.saturating_sub(vm.balloon_size());
memory_actual_size += vm.virtio_mem_plugged_size();
}
let device_tree = self
.vm
.as_ref()
.map(|vm| vm.device_tree().lock().unwrap().clone());
Ok(VmInfoResponse {
config: Box::new(config),
state,
memory_actual_size,
device_tree,
})
}
None => Err(VmError::VmNotCreated),
// In case of a migration, we emit the old VM info, as the VM is
// immutable during a migration.
if let VmOwnership::Migration {
vm_info_response, ..
} = &self.vm
{
return Ok(vm_info_response.clone());
}
let vm_config = self.vm_config.as_ref().ok_or(VmError::VmNotCreated)?;
let vm_config = vm_config.lock().unwrap().clone();
let state = match &self.vm {
VmOwnership::Owned(vm) => vm.get_state(),
VmOwnership::None => VmState::Created,
VmOwnership::Migration { .. } => unreachable!("migration path is handled above"),
};
let base_memory_actual_size =
vm_config.memory.total_size() - vm_config.memory.hotplugged_size();
let (memory_actual_size, device_tree) = match &self.vm {
VmOwnership::Owned(vm) => (
base_memory_actual_size.saturating_sub(vm.balloon_size())
+ vm.virtio_mem_plugged_size(),
Some(vm.device_tree().lock().unwrap().clone()),
),
VmOwnership::None => (base_memory_actual_size, None),
VmOwnership::Migration { .. } => unreachable!("migration path is handled above"),
};
Ok(VmInfoResponse {
config: Box::new(vm_config),
state,
memory_actual_size,
device_tree,
})
}
fn vmm_ping(&self) -> VmmPingResponse {
@@ -2241,6 +2331,7 @@ impl RequestHandler for Vmm {
// If a VM is booted, we first try to shut it down.
self.vm_shutdown()?;
}
VmOwnership::Migration { .. } => return Err(VmError::VmMigrating),
VmOwnership::None => {}
}
@@ -2268,6 +2359,7 @@ impl RequestHandler for Vmm {
VmOwnership::Owned(ref mut vm) => vm
.resize(desired_vcpus, desired_ram, desired_balloon)
.inspect_err(|e| error!("Error when resizing VM: {e:?}")),
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
if let Some(desired_vcpus) = desired_vcpus {
@@ -2292,6 +2384,7 @@ impl RequestHandler for Vmm {
match self.vm {
VmOwnership::Owned(ref mut vm) => vm.resize_disk(&id, desired_size),
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::ResizeDisk),
}
}
@@ -2305,6 +2398,7 @@ impl RequestHandler for Vmm {
.inspect_err(|e| error!("Error when resizing zone: {e:?}"))?;
Ok(())
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by setting the new desired ram.
let memory_config = &mut self.vm_config.as_ref().unwrap().lock().unwrap().memory;
@@ -2346,6 +2440,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2377,6 +2472,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2393,6 +2489,7 @@ impl RequestHandler for Vmm {
.inspect_err(|e| error!("Error when removing device from the VM: {e:?}"))?;
Ok(())
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
if let Some(ref config) = self.vm_config {
let mut config = config.lock().unwrap();
@@ -2427,6 +2524,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2455,6 +2553,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2491,6 +2590,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2519,6 +2619,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2547,6 +2648,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2575,6 +2677,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2608,6 +2711,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => {
// Update VmConfig by adding the new device.
let mut config = self.vm_config.as_ref().unwrap().lock().unwrap();
@@ -2627,6 +2731,7 @@ impl RequestHandler for Vmm {
.map(Some)
.map_err(VmError::SerializeJson)
}
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
@@ -2634,6 +2739,7 @@ impl RequestHandler for Vmm {
fn vm_power_button(&mut self) -> result::Result<(), VmError> {
match self.vm {
VmOwnership::Owned(ref mut vm) => vm.power_button(),
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
@@ -2641,6 +2747,7 @@ impl RequestHandler for Vmm {
fn vm_nmi(&mut self) -> result::Result<(), VmError> {
match self.vm {
VmOwnership::Owned(ref mut vm) => vm.nmi(),
VmOwnership::Migration { .. } => Err(VmError::VmMigrating),
VmOwnership::None => Err(VmError::VmNotRunning),
}
}
@@ -2655,6 +2762,11 @@ impl RequestHandler for Vmm {
"Can't receive a migration when a VM is already created"
)));
}
VmOwnership::Migration { .. } => {
return Err(MigratableError::MigrateReceive(anyhow!(
"There is already an ongoing migration"
)));
}
VmOwnership::None => {}
}
@@ -2725,10 +2837,27 @@ impl RequestHandler for Vmm {
Ok(())
}
/// Dispatches a migration.
///
/// Returns an error if the migration worker cannot be spawned. Once
/// spawned, [`Vmm::check_migration`] will be called after the thread exits
/// (on success, cancellation, or failure).
fn vm_send_migration(
&mut self,
send_data_migration: VmSendMigrationData,
) -> result::Result<(), MigratableError> {
match self.vm {
VmOwnership::Owned(_) => (),
VmOwnership::Migration { .. } => {
return Err(MigratableError::MigrateSend(anyhow!(
"There is already an ongoing migration"
)));
}
VmOwnership::None => {
return Err(MigratableError::MigrateSend(anyhow!("VM is not running")));
}
}
send_data_migration
.validate()
.context("Invalid send migration configuration")
@@ -2770,45 +2899,43 @@ impl RequestHandler for Vmm {
)));
}
event!("vm", "migration-started");
Self::send_migration(
vm,
#[cfg(all(feature = "kvm", target_arch = "x86_64"))]
self.hypervisor.as_ref(),
&send_data_migration,
initial_vm_state,
)
.map_err(|migration_err| {
error!("Migration failed: {migration_err:?}");
event!("vm", "migration-failed");
// Stop logging dirty pages only for non-local migrations
if !send_data_migration.local
&& let Err(e) = vm.stop_dirty_log()
{
return e;
}
// Only resume if the VM was originally running; a VM that was already
// paused before migration should remain paused after failure.
if initial_vm_state == VmState::Running
&& vm.get_state() == VmState::Paused
&& let Err(e) = vm.resume()
{
return e;
}
migration_err
let vm_info_snapshot = self.vm_info().map_err(|e| {
MigratableError::MigrateSend(anyhow!("Failed to query VM info snapshot: {e}"))
})?;
event!("vm", "migration-finished");
let check_migration_evt = self
.check_migration_evt
.try_clone()
.with_context(|| "Failed to clone check_migration_evt FD")
.map_err(MigratableError::MigrateSend)?;
// Shutdown the VM after the migration succeeded
self.exit_evt.write(1).map_err(|e| {
MigratableError::MigrateSend(anyhow!(
"Failed shutting down the VM after migration: {e:?}"
))
})
// Take VM ownership. This also means that API events can no longer
// change the VM (e.g. net device hotplug).
let vm = self
.vm
.take_owned_or(VmError::VmNotRunning)
.expect("should have VM ownership as we just checked it");
match MigrationWorker::spawn(
vm,
check_migration_evt,
send_data_migration,
#[cfg(all(feature = "kvm", target_arch = "x86_64"))]
self.hypervisor.clone(),
initial_vm_state,
) {
Ok(handle) => {
self.vm = VmOwnership::Migration {
migration_worker_handle: handle,
vm_info_response: vm_info_snapshot,
};
Ok(())
}
Err(e) => {
self.vm = VmOwnership::Owned(e.vm);
Err(MigratableError::MigrateSend(e.spawn_error.into()))
}
}
}
}

View File

@@ -191,6 +191,9 @@ pub enum Error {
#[error("VM is not running")]
VmNotRunning,
#[error("VM is currently migrating and can't be modified")]
VmMigrating,
#[error("Cannot clone EventFd")]
EventFdClone(#[source] io::Error),