diff --git a/block/src/qcow/metadata.rs b/block/src/qcow/metadata.rs index c78a2af0b..b4b64cabd 100644 --- a/block/src/qcow/metadata.rs +++ b/block/src/qcow/metadata.rs @@ -130,7 +130,7 @@ pub(crate) struct QcowState { } impl QcowMetadata { - pub(super) fn new(inner: QcowState) -> Self { + pub(crate) fn new(inner: QcowState) -> Self { QcowMetadata { inner: RwLock::new(inner), } @@ -238,6 +238,21 @@ impl QcowMetadata { Ok(()) } + /// Flushes dirty metadata caches and clears the dirty bit for + /// clean shutdown. + pub fn shutdown(&self) { + let mut inner = self.inner.write().unwrap(); + let _ = inner.sync_caches(); + let QcowState { + ref mut header, + ref mut raw_file, + .. + } = *inner; + if raw_file.file().is_writable() { + let _ = header.set_dirty_bit(raw_file.file_mut(), false); + } + } + /// Resizes the QCOW2 image to the given new size. Only grow is /// supported, shrink would require walking all L2 tables to reclaim /// clusters beyond the new size and risks data loss. @@ -260,6 +275,9 @@ impl QcowMetadata { cluster_size: u64, backing_file: Option<&dyn BackingRead>, ) -> io::Result> { + if address.checked_add(length as u64).is_none() { + return Ok(Vec::new()); + } let mut inner = self.inner.write().unwrap(); let mut actions = Vec::new(); @@ -476,6 +494,19 @@ impl QcowState { } } + /// Maps a single cluster region for a sequential read. + pub(crate) fn map_cluster_read( + &mut self, + address: u64, + count: usize, + has_backing_file: bool, + ) -> io::Result { + match self.try_map_read(address, count, has_backing_file)? { + Some(mapping) => Ok(mapping), + None => self.map_read_with_populate(address, count, has_backing_file), + } + } + /// Write path mapping. Always called under write lock. fn map_write( &mut self, diff --git a/block/src/qcow/mod.rs b/block/src/qcow/mod.rs index 06dce2c4c..c0b4e8c72 100644 --- a/block/src/qcow/mod.rs +++ b/block/src/qcow/mod.rs @@ -7,7 +7,7 @@ mod decoder; mod header; pub(crate) mod metadata; -mod qcow_raw_file; +pub(crate) mod qcow_raw_file; mod raw_file; mod refcount; mod util; @@ -37,6 +37,7 @@ use header::{ }; use libc::{EINVAL, EIO, ENOSPC}; use log::{error, warn}; +use metadata::ClusterReadMapping; use remain::sorted; use thiserror::Error; pub(crate) use util::MAX_NESTING_DEPTH; @@ -162,29 +163,22 @@ pub enum Error { pub type Result = std::result::Result; -trait BackingFileOps: Send + Seek + Read { - fn read_at(&mut self, address: u64, buf: &mut [u8]) -> std::io::Result<()> { - self.seek(SeekFrom::Start(address))?; - self.read_exact(buf) - } - fn clone_box(&self) -> Box; +/// Concrete backing file variants. +pub(crate) enum BackingKind { + /// Raw backing file. + Raw(RawFile), + /// QCOW2 backing parsed into metadata and raw file. + Qcow { + inner: Box, + backing: Option>, + }, + /// Full QcowFile used as backing, only in tests. + #[cfg(test)] + QcowFile(Box), } - -impl BackingFileOps for QcowFile { - fn clone_box(&self) -> Box { - Box::new(self.clone()) - } -} - -impl BackingFileOps for RawFile { - fn clone_box(&self) -> Box { - Box::new(self.clone()) - } -} - /// Backing file wrapper pub(crate) struct BackingFile { - inner: Box, + kind: BackingKind, virtual_size: u64, } @@ -217,56 +211,108 @@ impl BackingFile { None => detect_image_type(&mut raw_file)?, }; - let (inner, virtual_size): (Box, u64) = match backing_format { + let (kind, virtual_size) = match backing_format { ImageType::Raw => { let size = raw_file .seek(SeekFrom::End(0)) .map_err(Error::BackingFileIo)?; raw_file.rewind().map_err(Error::BackingFileIo)?; - (Box::new(raw_file), size) + (BackingKind::Raw(raw_file), size) } ImageType::Qcow2 => { - let backing_qcow = - QcowFile::from_with_nesting_depth(raw_file, max_nesting_depth - 1, sparse) + let (inner, nested_backing, _sparse) = + parse_qcow(raw_file, max_nesting_depth - 1, sparse) .map_err(|e| Error::BackingFileOpen(Box::new(e)))?; - let size = backing_qcow.virtual_size(); - (Box::new(backing_qcow), size) + let size = inner.header.size; + ( + BackingKind::Qcow { + inner: Box::new(inner), + backing: nested_backing.map(Box::new), + }, + size, + ) } }; - Ok(Some(Self { - inner, - virtual_size, - })) + Ok(Some(Self { kind, virtual_size })) + } + + /// Consume and return the kind and virtual size. + pub(crate) fn into_kind(self) -> (BackingKind, u64) { + (self.kind, self.virtual_size) } /// Read from backing file, returning zeros for any portion beyond backing file size. #[inline] - fn read_at(&mut self, address: u64, buf: &mut [u8]) -> std::io::Result<()> { + pub(crate) fn read_at(&mut self, address: u64, buf: &mut [u8]) -> std::io::Result<()> { if address >= self.virtual_size { - // Entire read is beyond backing file buf.fill(0); return Ok(()); } let available = (self.virtual_size - address) as usize; - if available >= buf.len() { - // Entire read is within backing file - self.inner.read_at(address, buf) + let (target, overflow) = if available >= buf.len() { + (buf, &mut [][..]) } else { - // Partial read, fill the rest with zeroes - self.inner.read_at(address, &mut buf[..available])?; - buf[available..].fill(0); - Ok(()) - } + buf.split_at_mut(available) + }; + Self::read_at_inner(&mut self.kind, address, target)?; + overflow.fill(0); + Ok(()) } -} -impl Clone for BackingFile { - fn clone(&self) -> Self { - Self { - inner: self.inner.clone_box(), - virtual_size: self.virtual_size, + fn read_at_inner(kind: &mut BackingKind, address: u64, buf: &mut [u8]) -> std::io::Result<()> { + match kind { + BackingKind::Raw(file) => { + file.seek(SeekFrom::Start(address))?; + file.read_exact(buf) + } + #[cfg(test)] + BackingKind::QcowFile(qcow) => { + qcow.seek(SeekFrom::Start(address))?; + qcow.read_exact(buf) + } + BackingKind::Qcow { inner, backing } => { + let has_backing = backing.is_some(); + let cluster_size = inner.raw_file.cluster_size(); + let mut pos = 0usize; + while pos < buf.len() { + let curr_addr = address + pos as u64; + let intra = inner.raw_file.cluster_offset(curr_addr) as usize; + let count = min(buf.len() - pos, cluster_size as usize - intra); + let mapping = inner.map_cluster_read(curr_addr, count, has_backing)?; + match mapping { + ClusterReadMapping::Zero { length } => { + buf[pos..pos + length as usize].fill(0); + } + ClusterReadMapping::Allocated { + offset: host_off, + length, + } => { + inner.raw_file.file_mut().seek(SeekFrom::Start(host_off))?; + inner + .raw_file + .file_mut() + .read_exact(&mut buf[pos..pos + length as usize])?; + } + ClusterReadMapping::Compressed { data } => { + buf[pos..pos + data.len()].copy_from_slice(&data); + } + ClusterReadMapping::Backing { + offset: backing_off, + length, + } => { + if let Some(bf) = backing.as_mut() { + bf.read_at(backing_off, &mut buf[pos..pos + length as usize])?; + } else { + buf[pos..pos + length as usize].fill(0); + } + } + } + pos += count; + } + Ok(()) + } } } } @@ -497,7 +543,7 @@ pub(crate) fn parse_qcow( /// # Ok(()) /// # } /// ``` -#[derive(Clone, Debug)] +#[derive(Debug)] pub struct QcowFile { raw_file: QcowRawFile, header: QcowHeader, @@ -605,11 +651,12 @@ impl QcowFile { Ok(qcow) } + #[cfg(test)] pub fn set_backing_file(&mut self, backing: Option>) { self.backing_file = backing.map(|b| { let virtual_size = b.virtual_size(); BackingFile { - inner: Box::new(*b), + kind: BackingKind::QcowFile(b), virtual_size, } }); diff --git a/block/src/qcow_sync.rs b/block/src/qcow_sync.rs index 2707f5dfb..22d361ada 100644 --- a/block/src/qcow_sync.rs +++ b/block/src/qcow_sync.rs @@ -2,80 +2,236 @@ // // SPDX-License-Identifier: Apache-2.0 AND BSD-3-Clause +use std::cmp::min; use std::collections::VecDeque; use std::fs::File; -use std::io::{self, Seek, SeekFrom}; -use std::os::fd::AsRawFd; -use std::sync::{Arc, Mutex}; +use std::os::fd::{AsRawFd, FromRawFd, OwnedFd, RawFd}; +use std::sync::Arc; +use std::{io, ptr, slice}; use vmm_sys_util::eventfd::EventFd; -use vmm_sys_util::write_zeroes::PunchHole; +use vmm_sys_util::write_zeroes::{PunchHole, WriteZeroesAt}; use crate::async_io::{ AsyncIo, AsyncIoError, AsyncIoResult, BorrowedDiskFd, DiskFile, DiskFileError, DiskFileResult, }; -use crate::qcow::{Error as QcowError, MAX_NESTING_DEPTH, QcowFile, RawFile, Result as QcowResult}; -use crate::{AsyncAdaptor, BlockBackend}; +use crate::qcow::metadata::{ + BackingRead, ClusterReadMapping, ClusterWriteMapping, DeallocAction, QcowMetadata, +}; +use crate::qcow::qcow_raw_file::QcowRawFile; +use crate::qcow::{ + BackingFile, BackingKind, Error as QcowError, MAX_NESTING_DEPTH, RawFile, Result as QcowResult, + parse_qcow, +}; + +/// Raw backing file using pread64 on a duplicated fd. +struct RawBacking { + fd: OwnedFd, + virtual_size: u64, +} + +// SAFETY: The only I/O operation is pread64 which is position independent +// and safe for concurrent use from multiple threads. +unsafe impl Sync for RawBacking {} + +impl BackingRead for RawBacking { + fn read_at(&self, address: u64, buf: &mut [u8]) -> io::Result<()> { + if address >= self.virtual_size { + buf.fill(0); + return Ok(()); + } + let available = (self.virtual_size - address) as usize; + if available >= buf.len() { + pread_exact(self.fd.as_raw_fd(), buf, address) + } else { + pread_exact(self.fd.as_raw_fd(), &mut buf[..available], address)?; + buf[available..].fill(0); + Ok(()) + } + } +} + +/// QCOW2 backing file with RwLock metadata and pread64 data reads. +/// +/// Read only because backing files never receive writes. Nested backing +/// files are handled recursively. +struct Qcow2MetadataBacking { + metadata: Arc, + data_fd: OwnedFd, + backing_file: Option>, +} + +// SAFETY: All reads go through QcowMetadata which uses RwLock +// and pread64 which is position independent and thread safe. +unsafe impl Sync for Qcow2MetadataBacking {} + +impl BackingRead for Qcow2MetadataBacking { + fn read_at(&self, address: u64, buf: &mut [u8]) -> io::Result<()> { + let virtual_size = self.metadata.virtual_size(); + if address >= virtual_size { + buf.fill(0); + return Ok(()); + } + let available = (virtual_size - address) as usize; + if available < buf.len() { + self.read_clusters(address, &mut buf[..available])?; + buf[available..].fill(0); + return Ok(()); + } + self.read_clusters(address, buf) + } +} + +impl Qcow2MetadataBacking { + /// Resolve cluster mappings via metadata then read allocated clusters + /// with pread64. + fn read_clusters(&self, address: u64, buf: &mut [u8]) -> io::Result<()> { + let total_len = buf.len(); + let has_backing = self.backing_file.is_some(); + + let mappings = self + .metadata + .map_clusters_for_read(address, total_len, has_backing)?; + + let mut buf_offset = 0usize; + for mapping in mappings { + match mapping { + ClusterReadMapping::Zero { length } => { + buf[buf_offset..buf_offset + length as usize].fill(0); + buf_offset += length as usize; + } + ClusterReadMapping::Allocated { + offset: host_offset, + length, + } => { + pread_exact( + self.data_fd.as_raw_fd(), + &mut buf[buf_offset..buf_offset + length as usize], + host_offset, + )?; + buf_offset += length as usize; + } + ClusterReadMapping::Compressed { data } => { + let len = data.len(); + buf[buf_offset..buf_offset + len].copy_from_slice(&data); + buf_offset += len; + } + ClusterReadMapping::Backing { + offset: backing_offset, + length, + } => { + self.backing_file.as_ref().unwrap().read_at( + backing_offset, + &mut buf[buf_offset..buf_offset + length as usize], + )?; + buf_offset += length as usize; + } + } + } + Ok(()) + } +} + +impl Drop for Qcow2MetadataBacking { + fn drop(&mut self) { + self.metadata.shutdown(); + } +} + +/// Construct a thread safe backing file reader. +fn shared_backing_from(bf: BackingFile) -> QcowResult> { + let (kind, virtual_size) = bf.into_kind(); + match kind { + BackingKind::Raw(raw_file) => { + // SAFETY: raw_file holds a valid open fd. + let dup_fd = unsafe { libc::dup(raw_file.as_raw_fd()) }; + if dup_fd < 0 { + return Err(QcowError::BackingFileIo(io::Error::last_os_error())); + } + // SAFETY: dup_fd is a freshly duplicated valid fd. + let fd = unsafe { OwnedFd::from_raw_fd(dup_fd) }; + Ok(Arc::new(RawBacking { fd, virtual_size })) + } + BackingKind::Qcow { inner, backing } => { + // SAFETY: inner.raw_file holds a valid open fd. + let dup_fd = unsafe { libc::dup(inner.raw_file.as_raw_fd()) }; + if dup_fd < 0 { + return Err(QcowError::BackingFileIo(io::Error::last_os_error())); + } + // SAFETY: dup_fd is a freshly duplicated valid fd. + let data_fd = unsafe { OwnedFd::from_raw_fd(dup_fd) }; + Ok(Arc::new(Qcow2MetadataBacking { + metadata: Arc::new(QcowMetadata::new(*inner)), + data_fd, + backing_file: backing.map(|bf| shared_backing_from(*bf)).transpose()?, + })) + } + #[cfg(test)] + BackingKind::QcowFile(_) => { + unreachable!("QcowFile variant is only used by set_backing_file() in tests") + } + } +} pub struct QcowDiskSync { - // FIXME: The Mutex serializes all QCOW2 I/O operations across queues, which - // is necessary for correctness but eliminates any parallelism benefit from - // multiqueue. QcowFile has internal mutable state (L2 cache, refcounts, file - // position) that is not safe to share across threads via Clone. - // - // A proper fix would require restructuring QcowFile to separate metadata - // operations (which need synchronization) from data I/O (which could be - // parallelized with per queue file descriptors). See #7560 for details. - qcow_file: Arc>, + metadata: Arc, + /// Shared across queues, resolved once at construction. + backing_file: Option>, + sparse: bool, + data_raw_file: QcowRawFile, } impl QcowDiskSync { pub fn new(file: File, direct_io: bool, backing_files: bool, sparse: bool) -> QcowResult { let max_nesting_depth = if backing_files { MAX_NESTING_DEPTH } else { 0 }; - let qcow_file = QcowFile::from_with_nesting_depth( - RawFile::new(file, direct_io), - max_nesting_depth, - sparse, - ) - .map_err(|e| match e { - QcowError::MaxNestingDepthExceeded if !backing_files => QcowError::BackingFilesDisabled, - other => other, - })?; + let (inner, backing_file, sparse) = + parse_qcow(RawFile::new(file, direct_io), max_nesting_depth, sparse).map_err(|e| { + match e { + QcowError::MaxNestingDepthExceeded if !backing_files => { + QcowError::BackingFilesDisabled + } + other => other, + } + })?; + let data_raw_file = inner.raw_file.clone(); Ok(QcowDiskSync { - qcow_file: Arc::new(Mutex::new(qcow_file)), + metadata: Arc::new(QcowMetadata::new(inner)), + backing_file: backing_file.map(shared_backing_from).transpose()?, + sparse, + data_raw_file, }) } } impl DiskFile for QcowDiskSync { fn logical_size(&mut self) -> DiskFileResult { - self.qcow_file - .lock() - .unwrap() - .seek(SeekFrom::End(0)) - .map_err(DiskFileError::Size) + Ok(self.metadata.virtual_size()) } fn physical_size(&mut self) -> DiskFileResult { - self.qcow_file.lock().unwrap().physical_size().map_err(|e| { - let io_inner = match e { - crate::Error::GetFileMetadata(e) => e, - _ => unreachable!(), - }; - DiskFileError::Size(io_inner) - }) + self.data_raw_file + .physical_size() + .map_err(DiskFileError::Size) } fn new_async_io(&self, _ring_depth: u32) -> DiskFileResult> { - Ok(Box::new(QcowSync::new(Arc::clone(&self.qcow_file))) as Box) + Ok(Box::new(QcowSync::new( + Arc::clone(&self.metadata), + self.data_raw_file.clone(), + self.backing_file.as_ref().map(Arc::clone), + self.sparse, + )) as Box) } fn resize(&mut self, size: u64) -> DiskFileResult<()> { - self.qcow_file - .lock() - .unwrap() + if self.backing_file.is_some() { + return Err(DiskFileError::ResizeError(io::Error::other( + "resize not supported with backing file", + ))); + } + self.metadata .resize(size) - .map_err(|e| DiskFileError::ResizeError(io::Error::other(e))) + .map_err(DiskFileError::ResizeError) } fn supports_sparse_operations(&self) -> bool { @@ -87,20 +243,38 @@ impl DiskFile for QcowDiskSync { } fn fd(&mut self) -> BorrowedDiskFd<'_> { - BorrowedDiskFd::new(self.qcow_file.lock().unwrap().as_raw_fd()) + BorrowedDiskFd::new(self.data_raw_file.as_raw_fd()) + } +} + +impl Drop for QcowDiskSync { + fn drop(&mut self) { + self.metadata.shutdown(); } } pub struct QcowSync { - qcow_file: Arc>, + metadata: Arc, + data_file: QcowRawFile, + /// See the backing_file field on QcowDiskSync. + backing_file: Option>, + sparse: bool, eventfd: EventFd, completion_list: VecDeque<(u64, i32)>, } impl QcowSync { - pub fn new(qcow_file: Arc>) -> Self { + fn new( + metadata: Arc, + data_file: QcowRawFile, + backing_file: Option>, + sparse: bool, + ) -> Self { QcowSync { - qcow_file, + metadata, + data_file, + backing_file, + sparse, eventfd: EventFd::new(libc::EFD_NONBLOCK) .expect("Failed creating EventFd for QcowSync"), completion_list: VecDeque::new(), @@ -108,7 +282,152 @@ impl QcowSync { } } -impl AsyncAdaptor for QcowFile {} +// -- Position independent I/O helpers -- +// +// Duplicated file descriptors share the kernel file description and thus the +// file position. Using seek then read from multiple queues races on that +// shared position. pread64 and pwrite64 are atomic and never touch the position. + +/// Read exactly the requested bytes at offset, looping on short reads. +fn pread_exact(fd: RawFd, buf: &mut [u8], offset: u64) -> io::Result<()> { + let mut total = 0usize; + while total < buf.len() { + // SAFETY: buf and fd are valid for the lifetime of the call. + let ret = unsafe { + libc::pread64( + fd, + buf[total..].as_mut_ptr() as *mut libc::c_void, + buf.len() - total, + (offset + total as u64) as libc::off_t, + ) + }; + if ret < 0 { + return Err(io::Error::last_os_error()); + } + if ret == 0 { + return Err(io::Error::from(io::ErrorKind::UnexpectedEof)); + } + total += ret as usize; + } + Ok(()) +} + +/// Write all bytes to fd at offset, looping on short writes. +fn pwrite_all(fd: RawFd, buf: &[u8], offset: u64) -> io::Result<()> { + let mut total = 0usize; + while total < buf.len() { + // SAFETY: buf and fd are valid for the lifetime of the call. + let ret = unsafe { + libc::pwrite64( + fd, + buf[total..].as_ptr() as *const libc::c_void, + buf.len() - total, + (offset + total as u64) as libc::off_t, + ) + }; + if ret < 0 { + return Err(io::Error::last_os_error()); + } + if ret == 0 { + return Err(io::Error::other("pwrite64 wrote 0 bytes")); + } + total += ret as usize; + } + Ok(()) +} + +// -- iovec helper functions -- +// +// Operate on the iovec array as a flat byte stream. + +/// Copy data into iovecs starting at the given byte offset. +/// +/// # Safety +/// Caller must ensure iovecs point to valid, writable memory of sufficient size. +unsafe fn scatter_to_iovecs(iovecs: &[libc::iovec], start: usize, data: &[u8]) { + let mut remaining = data; + let mut pos = 0usize; + for iov in iovecs { + let iov_end = pos + iov.iov_len; + if iov_end <= start || remaining.is_empty() { + pos = iov_end; + continue; + } + let iov_start = start.saturating_sub(pos); + let available = iov.iov_len - iov_start; + let count = min(available, remaining.len()); + // SAFETY: iov_base is valid for iov_len bytes per caller contract. + unsafe { + let dst = (iov.iov_base as *mut u8).add(iov_start); + ptr::copy_nonoverlapping(remaining.as_ptr(), dst, count); + } + remaining = &remaining[count..]; + if remaining.is_empty() { + break; + } + pos = iov_end; + } +} + +/// Zero fill iovecs starting at the given byte offset for the given length. +/// +/// # Safety +/// Caller must ensure iovecs point to valid, writable memory of sufficient size. +unsafe fn zero_fill_iovecs(iovecs: &[libc::iovec], start: usize, len: usize) { + let mut remaining = len; + let mut pos = 0usize; + for iov in iovecs { + let iov_end = pos + iov.iov_len; + if iov_end <= start || remaining == 0 { + pos = iov_end; + continue; + } + let iov_start = start.saturating_sub(pos); + let available = iov.iov_len - iov_start; + let count = min(available, remaining); + // SAFETY: iov_base is valid for iov_len bytes per caller contract. + unsafe { + let dst = (iov.iov_base as *mut u8).add(iov_start); + ptr::write_bytes(dst, 0, count); + } + remaining -= count; + if remaining == 0 { + break; + } + pos = iov_end; + } +} + +/// Gather bytes from iovecs starting at the given byte offset into a Vec. +/// +/// # Safety +/// Caller must ensure iovecs point to valid, readable memory of sufficient size. +unsafe fn gather_from_iovecs(iovecs: &[libc::iovec], start: usize, len: usize) -> Vec { + let mut result = Vec::with_capacity(len); + let mut remaining = len; + let mut pos = 0usize; + for iov in iovecs { + let iov_end = pos + iov.iov_len; + if iov_end <= start || remaining == 0 { + pos = iov_end; + continue; + } + let iov_start = start.saturating_sub(pos); + let available = iov.iov_len - iov_start; + let count = min(available, remaining); + // SAFETY: iov_base is valid for iov_len bytes per caller contract. + unsafe { + let src = (iov.iov_base as *const u8).add(iov_start); + result.extend_from_slice(slice::from_raw_parts(src, count)); + } + remaining -= count; + if remaining == 0 { + break; + } + pos = iov_end; + } + result +} impl AsyncIo for QcowSync { fn notifier(&self) -> &EventFd { @@ -121,13 +440,61 @@ impl AsyncIo for QcowSync { iovecs: &[libc::iovec], user_data: u64, ) -> AsyncIoResult<()> { - self.qcow_file.lock().unwrap().read_vectored_sync( - offset, - iovecs, - user_data, - &self.eventfd, - &mut self.completion_list, - ) + let address = offset as u64; + let total_len: usize = iovecs.iter().map(|v| v.iov_len).sum(); + + let has_backing = self.backing_file.is_some(); + let mappings = self + .metadata + .map_clusters_for_read(address, total_len, has_backing) + .map_err(AsyncIoError::ReadVectored)?; + + let mut buf_offset = 0usize; + for mapping in mappings { + match mapping { + ClusterReadMapping::Zero { length } => { + // SAFETY: iovecs point to valid guest memory buffers + unsafe { zero_fill_iovecs(iovecs, buf_offset, length as usize) }; + buf_offset += length as usize; + } + ClusterReadMapping::Allocated { + offset: host_offset, + length, + } => { + let mut buf = vec![0u8; length as usize]; + pread_exact(self.data_file.as_raw_fd(), &mut buf, host_offset) + .map_err(AsyncIoError::ReadVectored)?; + // SAFETY: iovecs point to valid guest memory buffers + unsafe { scatter_to_iovecs(iovecs, buf_offset, &buf) }; + buf_offset += length as usize; + } + ClusterReadMapping::Compressed { data } => { + let len = data.len(); + // SAFETY: iovecs point to valid guest memory buffers + unsafe { scatter_to_iovecs(iovecs, buf_offset, &data) }; + buf_offset += len; + } + ClusterReadMapping::Backing { + offset: backing_offset, + length, + } => { + let mut buf = vec![0u8; length as usize]; + self.backing_file + .as_ref() + .unwrap() + .read_at(backing_offset, &mut buf) + .map_err(AsyncIoError::ReadVectored)?; + // SAFETY: iovecs point to valid guest memory buffers + unsafe { scatter_to_iovecs(iovecs, buf_offset, &buf) }; + buf_offset += length as usize; + } + } + } + + self.completion_list + .push_back((user_data, total_len as i32)); + self.eventfd.write(1).unwrap(); + Ok(()) } fn write_vectored( @@ -136,21 +503,65 @@ impl AsyncIo for QcowSync { iovecs: &[libc::iovec], user_data: u64, ) -> AsyncIoResult<()> { - self.qcow_file.lock().unwrap().write_vectored_sync( - offset, - iovecs, - user_data, - &self.eventfd, - &mut self.completion_list, - ) + let address = offset as u64; + let total_len: usize = iovecs.iter().map(|v| v.iov_len).sum(); + let mut buf_offset = 0usize; + + while buf_offset < total_len { + let curr_addr = address + buf_offset as u64; + let cluster_size = self.metadata.cluster_size(); + let intra_offset = self.metadata.cluster_offset(curr_addr); + let remaining_in_cluster = (cluster_size - intra_offset) as usize; + let count = min(total_len - buf_offset, remaining_in_cluster); + + // Read backing data for COW if this is a partial cluster + // write to an unallocated cluster with a backing file. + let backing_data = if let Some(backing) = self + .backing_file + .as_ref() + .filter(|_| intra_offset != 0 || count < cluster_size as usize) + { + let cluster_begin = curr_addr - intra_offset; + let mut data = vec![0u8; cluster_size as usize]; + backing + .read_at(cluster_begin, &mut data) + .map_err(AsyncIoError::WriteVectored)?; + Some(data) + } else { + None + }; + + let mapping = self + .metadata + .map_cluster_for_write(curr_addr, backing_data) + .map_err(AsyncIoError::WriteVectored)?; + + match mapping { + ClusterWriteMapping::Allocated { + offset: host_offset, + } => { + // SAFETY: iovecs point to valid guest memory buffers + let buf = unsafe { gather_from_iovecs(iovecs, buf_offset, count) }; + pwrite_all(self.data_file.as_raw_fd(), &buf, host_offset) + .map_err(AsyncIoError::WriteVectored)?; + } + } + buf_offset += count; + } + + self.completion_list + .push_back((user_data, total_len as i32)); + self.eventfd.write(1).unwrap(); + Ok(()) } fn fsync(&mut self, user_data: Option) -> AsyncIoResult<()> { - self.qcow_file.lock().unwrap().fsync_sync( - user_data, - &self.eventfd, - &mut self.completion_list, - ) + self.metadata.flush().map_err(AsyncIoError::Fsync)?; + if let Some(user_data) = user_data { + self.completion_list.push_back((user_data, 0)); + self.eventfd.write(1).unwrap(); + } + Ok(()) } fn next_completed_request(&mut self) -> Option<(u64, i32)> { @@ -158,26 +569,49 @@ impl AsyncIo for QcowSync { } fn punch_hole(&mut self, offset: u64, length: u64, user_data: u64) -> AsyncIoResult<()> { - // For QCOW2, punch_hole calls deallocate_cluster + let virtual_size = self.metadata.virtual_size(); + let cluster_size = self.metadata.cluster_size(); + let result = self - .qcow_file - .lock() - .unwrap() - .punch_hole(offset, length) - .map(|_| 0i32) + .metadata + .deallocate_bytes( + offset, + length as usize, + self.sparse, + virtual_size, + cluster_size, + self.backing_file.as_deref(), + ) .map_err(AsyncIoError::PunchHole); match result { - Ok(res) => { - self.completion_list.push_back((user_data, res)); + Ok(actions) => { + for action in actions { + match action { + DeallocAction::PunchHole { + host_offset, + length, + } => { + let _ = self.data_file.file_mut().punch_hole(host_offset, length); + } + DeallocAction::WriteZeroes { + host_offset, + length, + } => { + let _ = self + .data_file + .file_mut() + .write_zeroes_at(host_offset, length); + } + } + } + self.completion_list.push_back((user_data, 0)); self.eventfd.write(1).unwrap(); Ok(()) } Err(e) => { - // CRITICAL: Always signal completion even on error to avoid hangs - let errno = if let AsyncIoError::PunchHole(io_err) = &e { - let err = io_err.raw_os_error().unwrap_or(libc::EIO); - -err + let errno = if let AsyncIoError::PunchHole(ref io_err) = e { + -io_err.raw_os_error().unwrap_or(libc::EIO) } else { -libc::EIO }; @@ -189,85 +623,90 @@ impl AsyncIo for QcowSync { } fn write_zeroes(&mut self, offset: u64, length: u64, user_data: u64) -> AsyncIoResult<()> { - // For QCOW2, write_zeroes is implemented by deallocating clusters via punch_hole. - // This is more efficient than writing actual zeros and reduces disk usage. + // For QCOW2 write_zeroes uses cluster deallocation, same as punch_hole. // Unallocated clusters inherently read as zero in the QCOW2 format. - let result = self - .qcow_file - .lock() - .unwrap() - .punch_hole(offset, length) - .map(|_| 0i32) - .map_err(AsyncIoError::WriteZeroes); - - match result { - Ok(res) => { - self.completion_list.push_back((user_data, res)); - self.eventfd.write(1).unwrap(); - Ok(()) - } - Err(e) => { - // Always signal completion even on error to avoid hangs - let errno = if let AsyncIoError::WriteZeroes(io_err) = &e { - let err = io_err.raw_os_error().unwrap_or(libc::EIO); - -err - } else { - -libc::EIO - }; - self.completion_list.push_back((user_data, errno)); - self.eventfd.write(1).unwrap(); - Ok(()) - } - } + self.punch_hole(offset, length, user_data) } } #[cfg(test)] mod unit_tests { - use std::io::{Read, Seek, SeekFrom, Write}; + use std::io::{Seek, SeekFrom, Write}; use vmm_sys_util::tempfile::TempFile; use super::*; - use crate::qcow::{QcowFile, QcowHeader, RawFile}; + use crate::async_io::DiskFile; + use crate::qcow::{QcowFile, RawFile}; + + fn create_disk_with_data( + file_size: u64, + data: &[u8], + offset: u64, + sparse: bool, + ) -> (TempFile, QcowDiskSync) { + let temp_file = TempFile::new().unwrap(); + { + let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false); + let mut qcow_file = QcowFile::new(raw_file, 3, file_size, sparse).unwrap(); + qcow_file.seek(SeekFrom::Start(offset)).unwrap(); + qcow_file.write_all(data).unwrap(); + qcow_file.flush().unwrap(); + } + let disk = QcowDiskSync::new( + temp_file.as_file().try_clone().unwrap(), + false, + false, + sparse, + ) + .unwrap(); + (temp_file, disk) + } + + fn async_read(disk: &QcowDiskSync, offset: u64, len: usize) -> Vec { + let mut async_io = disk.new_async_io(1).unwrap(); + let mut buf = vec![0xFFu8; len]; + let iovec = libc::iovec { + iov_base: buf.as_mut_ptr() as *mut libc::c_void, + iov_len: buf.len(), + }; + async_io + .read_vectored(offset as libc::off_t, &[iovec], 1) + .unwrap(); + let (user_data, result) = async_io.next_completed_request().unwrap(); + assert_eq!(user_data, 1); + assert_eq!(result as usize, len, "read should return requested length"); + buf + } + + fn async_write(disk: &QcowDiskSync, offset: u64, data: &[u8]) { + let mut async_io = disk.new_async_io(1).unwrap(); + let iovec = libc::iovec { + iov_base: data.as_ptr() as *mut libc::c_void, + iov_len: data.len(), + }; + async_io + .write_vectored(offset as libc::off_t, &[iovec], 1) + .unwrap(); + let (user_data, result) = async_io.next_completed_request().unwrap(); + assert_eq!(user_data, 1); + assert_eq!(result as usize, data.len()); + } #[test] fn test_qcow_async_punch_hole_completion() { - // Create a QCOW2 image with valid header - let temp_file = TempFile::new().unwrap(); - let raw_file = RawFile::new(temp_file.into_file(), false); - let file_size = 1024 * 1024 * 100; // 100MB - let mut qcow_file = QcowFile::new(raw_file, 3, file_size, true).unwrap(); + let data = vec![0xDD; 128 * 1024]; + let offset = 0u64; + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true); - // Write some data - let data = vec![0xDD; 128 * 1024]; // 128KB - let offset = 0; - qcow_file.seek(SeekFrom::Start(offset)).unwrap(); - qcow_file.write_all(&data).unwrap(); - qcow_file.flush().unwrap(); - - // Create async wrapper - let qcow_file = Arc::new(Mutex::new(qcow_file)); - let mut async_qcow = QcowSync::new(qcow_file.clone()); - - // Punch hole - async_qcow - .punch_hole(offset, data.len() as u64, 100) - .unwrap(); - - // Verify completion event was generated - let (user_data, result) = async_qcow.next_completed_request().unwrap(); + let mut async_io = disk.new_async_io(1).unwrap(); + async_io.punch_hole(offset, data.len() as u64, 100).unwrap(); + let (user_data, result) = async_io.next_completed_request().unwrap(); assert_eq!(user_data, 100); assert_eq!(result, 0, "punch_hole should succeed"); + drop(async_io); - // Verify data reads as zeros - let mut read_buf = vec![0; data.len()]; - qcow_file - .lock() - .unwrap() - .seek(SeekFrom::Start(offset)) - .unwrap(); - qcow_file.lock().unwrap().read_exact(&mut read_buf).unwrap(); + let read_buf = async_read(&disk, offset, data.len()); assert!( read_buf.iter().all(|&b| b == 0), "Punched hole should read as zeros" @@ -276,41 +715,20 @@ mod unit_tests { #[test] fn test_qcow_async_write_zeroes_completion() { - // Create a QCOW2 image with valid header - let temp_file = TempFile::new().unwrap(); - let raw_file = RawFile::new(temp_file.into_file(), false); - let file_size = 1024 * 1024 * 100; // 100MB - let mut qcow_file = QcowFile::new(raw_file, 3, file_size, true).unwrap(); + let data = vec![0xEE; 256 * 1024]; + let offset = 64 * 1024u64; + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true); - // Write some data - let data = vec![0xEE; 256 * 1024]; // 256KB - let offset = 64 * 1024; // Start at 64KB offset - qcow_file.seek(SeekFrom::Start(offset)).unwrap(); - qcow_file.write_all(&data).unwrap(); - qcow_file.flush().unwrap(); - - // Create async wrapper - let qcow_file = Arc::new(Mutex::new(qcow_file)); - let mut async_qcow = QcowSync::new(qcow_file.clone()); - - // Write zeros - async_qcow + let mut async_io = disk.new_async_io(1).unwrap(); + async_io .write_zeroes(offset, data.len() as u64, 200) .unwrap(); - - // Verify completion event was generated - let (user_data, result) = async_qcow.next_completed_request().unwrap(); + let (user_data, result) = async_io.next_completed_request().unwrap(); assert_eq!(user_data, 200); assert_eq!(result, 0, "write_zeroes should succeed"); + drop(async_io); - // Verify data reads as zeros - let mut read_buf = vec![0; data.len()]; - qcow_file - .lock() - .unwrap() - .seek(SeekFrom::Start(offset)) - .unwrap(); - qcow_file.lock().unwrap().read_exact(&mut read_buf).unwrap(); + let read_buf = async_read(&disk, offset, data.len()); assert!( read_buf.iter().all(|&b| b == 0), "Zeroed region should read as zeros" @@ -319,164 +737,84 @@ mod unit_tests { #[test] fn test_qcow_async_multiple_operations() { - // Create a QCOW2 image with valid header - let temp_file = TempFile::new().unwrap(); - let raw_file = RawFile::new(temp_file.into_file(), false); - let file_size = 1024 * 1024 * 100; // 100MB - let mut qcow_file = QcowFile::new(raw_file, 3, file_size, true).unwrap(); + let data = vec![0xFF; 64 * 1024]; + let (_temp, _) = create_disk_with_data(100 * 1024 * 1024, &[], 0, true); - // Write data at multiple offsets - let data = vec![0xFF; 64 * 1024]; // 64KB chunks - for i in 0..4 { - let offset = i * 128 * 1024; // 128KB spacing - qcow_file.seek(SeekFrom::Start(offset)).unwrap(); - qcow_file.write_all(&data).unwrap(); + // Write data at multiple offsets via QcowFile first, then punch + { + let temp_file = _temp.as_file().try_clone().unwrap(); + let raw_file = RawFile::new(temp_file, false); + let mut qcow_file = QcowFile::from(raw_file).unwrap(); + for i in 0..4u64 { + let off = i * 128 * 1024; + qcow_file.seek(SeekFrom::Start(off)).unwrap(); + qcow_file.write_all(&data).unwrap(); + } + qcow_file.flush().unwrap(); } - qcow_file.flush().unwrap(); - // Create async wrapper - let qcow_file = Arc::new(Mutex::new(qcow_file)); - let mut async_qcow = QcowSync::new(qcow_file.clone()); + let disk = + QcowDiskSync::new(_temp.as_file().try_clone().unwrap(), false, false, true).unwrap(); - // Queue multiple punch_hole operations - async_qcow.punch_hole(0, 64 * 1024, 1).unwrap(); - async_qcow.punch_hole(128 * 1024, 64 * 1024, 2).unwrap(); - async_qcow.punch_hole(256 * 1024, 64 * 1024, 3).unwrap(); + let mut async_io = disk.new_async_io(1).unwrap(); - // Verify all completions - let (user_data, result) = async_qcow.next_completed_request().unwrap(); - assert_eq!(user_data, 1); - assert_eq!(result, 0); + async_io.punch_hole(0, 64 * 1024, 1).unwrap(); + async_io.punch_hole(128 * 1024, 64 * 1024, 2).unwrap(); + async_io.punch_hole(256 * 1024, 64 * 1024, 3).unwrap(); - let (user_data, result) = async_qcow.next_completed_request().unwrap(); - assert_eq!(user_data, 2); - assert_eq!(result, 0); - - let (user_data, result) = async_qcow.next_completed_request().unwrap(); - assert_eq!(user_data, 3); - assert_eq!(result, 0); - - // Verify no more completions - assert!(async_qcow.next_completed_request().is_none()); + let (ud, res) = async_io.next_completed_request().unwrap(); + assert_eq!(ud, 1); + assert_eq!(res, 0); + let (ud, res) = async_io.next_completed_request().unwrap(); + assert_eq!(ud, 2); + assert_eq!(res, 0); + let (ud, res) = async_io.next_completed_request().unwrap(); + assert_eq!(ud, 3); + assert_eq!(res, 0); + assert!(async_io.next_completed_request().is_none()); } #[test] - fn test_qcow_punch_hole_with_shared_instance() { - // This test verifies that with Arc>, multiple async I/O operations - // share the same QcowFile instance and see each other's changes. + fn test_qcow_punch_hole_then_read() { + // Verify that after punch_hole, a second async_io sees zeros. + let data = vec![0xAB; 128 * 1024]; + let offset = 0u64; + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true); - // Create a QCOW2 image - let temp_file = TempFile::new().unwrap(); - let raw_file = RawFile::new(temp_file.into_file(), false); - let file_size = 1024 * 1024 * 100; // 100MB - let mut qcow_file = QcowFile::new(raw_file, 3, file_size, true).unwrap(); - - // Write some data at offset 0 - let data = vec![0xAB; 128 * 1024]; // 128KB of 0xAB pattern - let offset = 0; - qcow_file.seek(SeekFrom::Start(offset)).unwrap(); - qcow_file.write_all(&data).unwrap(); - qcow_file.flush().unwrap(); - - let qcow_shared = Arc::new(Mutex::new(qcow_file)); - - // First async I/O: punch hole - let mut async_qcow1 = QcowSync::new(qcow_shared.clone()); - async_qcow1 + let mut async_io1 = disk.new_async_io(1).unwrap(); + async_io1 .punch_hole(offset, data.len() as u64, 100) .unwrap(); - - // Verify punch_hole completed - let (user_data, result) = async_qcow1.next_completed_request().unwrap(); + let (user_data, result) = async_io1.next_completed_request().unwrap(); assert_eq!(user_data, 100); - assert_eq!(result, 0, "punch_hole should succeed"); + assert_eq!(result, 0); + drop(async_io1); - // Second async I/O: read from same shared instance - // This should see the deallocated cluster because they share the same QcowFile - let mut read_buf = vec![0xFF; data.len()]; - qcow_shared - .lock() - .unwrap() - .seek(SeekFrom::Start(offset)) - .unwrap(); - qcow_shared - .lock() - .unwrap() - .read_exact(&mut read_buf) - .unwrap(); - - // The read should return zeros because the cluster was deallocated + // Read via second async_io, should see zeros + let read_buf = async_read(&disk, offset, data.len()); assert!( read_buf.iter().all(|&b| b == 0), - "After punch_hole, shared QcowFile instance should read zeros from deallocated cluster" + "After punch_hole, read should return zeros" ); } #[test] fn test_qcow_disk_sync_punch_hole_with_new_async_io() { - // This test simulates the EXACT real usage pattern: QcowDiskSync.new_async_io() - // creates a new QcowSync with a cloned QcowFile for each I/O operation. + // Simulates the real usage pattern of write data, punch hole, then read back. + let data = vec![0xCD; 64 * 1024]; // one cluster + let offset = 1024 * 1024u64; // 1MB offset + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true); - use std::io::Write; - - use crate::async_io::DiskFile; - - // Create a QCOW2 image - let temp_file = TempFile::new().unwrap(); - let file_size = 1024 * 1024 * 100; // 100MB - - { - let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false); - let mut qcow_file = QcowFile::new(raw_file, 3, file_size, true).unwrap(); - - // Write data at offset 1MB - use single cluster (64KB) to simplify test - let data = vec![0xCD; 64 * 1024]; // 64KB (one cluster) - let offset = 1024 * 1024u64; - qcow_file.seek(SeekFrom::Start(offset)).unwrap(); - qcow_file.write_all(&data).unwrap(); - qcow_file.flush().unwrap(); - } - - // Open with QcowDiskSync (like real code does) - let disk = - QcowDiskSync::new(temp_file.as_file().try_clone().unwrap(), false, true, true).unwrap(); - - // First async I/O: punch hole (simulates DISCARD command) + // Punch hole to simulate DISCARD let mut async_io1 = disk.new_async_io(1).unwrap(); - let offset = 1024 * 1024u64; - let length = 64 * 1024u64; // Single cluster - async_io1.punch_hole(offset, length, 1).unwrap(); + async_io1.punch_hole(offset, data.len() as u64, 1).unwrap(); let (user_data, result) = async_io1.next_completed_request().unwrap(); assert_eq!(user_data, 1); assert_eq!(result, 0, "punch_hole should succeed"); drop(async_io1); - // Second async I/O: read from the same location (simulates READ command) - let mut async_io2 = disk.new_async_io(1).unwrap(); - let mut read_buf = vec![0xFF; length as usize]; - let iovec = libc::iovec { - iov_base: read_buf.as_mut_ptr() as *mut libc::c_void, - iov_len: read_buf.len(), - }; - - // These assertions are critical to prevent compiler optimization bugs - // that can reorder operations. Without them, the test can fail even - // though the QCOW2 implementation is correct. - assert_eq!(iovec.iov_base as *const u8, read_buf.as_ptr()); - assert_eq!(iovec.iov_len, read_buf.len()); - - async_io2 - .read_vectored(offset as libc::off_t, &[iovec], 2) - .unwrap(); - - let (user_data, result) = async_io2.next_completed_request().unwrap(); - assert_eq!(user_data, 2); - assert_eq!( - result as usize, length as usize, - "read should complete successfully" - ); - - // Verify the data is all zeros + // Read from the same location to verify + let read_buf = async_read(&disk, offset, data.len()); assert!( read_buf.iter().all(|&b| b == 0), "After punch_hole via new_async_io, read should return zeros" @@ -484,21 +822,54 @@ mod unit_tests { } #[test] - fn backing_files_disabled_error() { - let header = - QcowHeader::create_for_size_and_path(3, 0x10_0000, Some("/path/to/backing/file")) - .expect("Failed to create header."); - let temp_file = TempFile::new().unwrap(); - let mut raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false); - header - .write_to(&mut raw_file) - .expect("Failed to write header."); + fn test_qcow_async_read_write_roundtrip() { + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &[], 0, true); - let file = temp_file.into_file(); - match QcowDiskSync::new(file, false, false, true) { - Err(QcowError::BackingFilesDisabled) => {} - Err(other) => panic!("Expected BackingFilesDisabled, got: {other:?}"), - Ok(_) => panic!("Expected BackingFilesDisabled error, but succeeded"), - } + let data = vec![0x42u8; 64 * 1024]; + let offset = 0u64; + + async_write(&disk, offset, &data); + + let mut async_io = disk.new_async_io(1).unwrap(); + async_io.fsync(Some(10)).unwrap(); + let (ud, res) = async_io.next_completed_request().unwrap(); + assert_eq!(ud, 10); + assert_eq!(res, 0); + drop(async_io); + + let read_buf = async_read(&disk, offset, data.len()); + assert_eq!(read_buf, data, "Read-back should match written data"); + } + + #[test] + fn test_qcow_async_read_unallocated() { + // Reading from an unallocated region should return zeros. + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &[], 0, true); + let read_buf = async_read(&disk, 0, 64 * 1024); + assert!( + read_buf.iter().all(|&b| b == 0), + "Unallocated region should read as zeros" + ); + } + + #[test] + fn test_qcow_async_cross_cluster_read_write() { + let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &[], 0, true); + + // Default cluster size is 64KB. Write 96KB starting at 32KB to cross the boundary. + let data: Vec = (0..96 * 1024).map(|i| (i % 251) as u8).collect(); + let offset = 32 * 1024u64; + + async_write(&disk, offset, &data); + + let mut async_io = disk.new_async_io(1).unwrap(); + async_io.fsync(Some(99)).unwrap(); + drop(async_io); + + let read_buf = async_read(&disk, offset, data.len()); + assert_eq!( + read_buf, data, + "Cross-cluster read should match written data" + ); } }