Files
cloud-hypervisor/block/src/qcow_async.rs
Julian Schindel 2b6e9df4e3 block: replace as <pointer> casts with safer alternatives
`as` casts can change mutability, which quickly leads to undefined
behavior.

Signed-off-by: Julian Schindel <mail@arctic-alpaca.de>
2026-05-03 08:38:31 +00:00

1087 lines
37 KiB
Rust

// Copyright © 2021 Intel Corporation
//
// Copyright 2026 The Cloud Hypervisor Authors. All rights reserved.
//
// SPDX-License-Identifier: Apache-2.0 AND BSD-3-Clause
//! QCOW2 async disk backend.
use std::cmp::{max, min};
use std::collections::VecDeque;
use std::io;
use std::os::unix::io::AsRawFd;
use std::sync::Arc;
use io_uring::{IoUring, opcode, types};
use vmm_sys_util::eventfd::EventFd;
use vmm_sys_util::write_zeroes::{PunchHole, WriteZeroesAt};
use crate::async_io::{AsyncIo, AsyncIoError, AsyncIoResult};
use crate::qcow::decoder::Decoder;
use crate::qcow::metadata::{
BackingRead, ClusterReadMapping, ClusterWriteMapping, DeallocAction, QcowMetadata,
};
use crate::qcow::qcow_raw_file::QcowRawFile;
use crate::qcow_common::{
AlignedBuf, aligned_pread, aligned_pwrite, decompress_cluster, gather_from_iovecs_into,
pread_alloc, pread_exact, pwrite_all, scatter_to_iovecs, zero_fill_iovecs,
};
use crate::{BatchRequest, RequestType, SECTOR_SIZE};
/// Per queue QCOW2 I/O worker using io_uring.
///
/// Reads against fully allocated single mapping clusters are submitted
/// to io_uring for true asynchronous completion. All other cluster
/// types (zero, compressed, backing) and multi mapping reads fall back
/// to synchronous I/O with synthetic completions.
///
/// Writes are synchronous because metadata allocation must complete
/// before the host offset is known.
pub struct QcowAsync {
metadata: Arc<QcowMetadata>,
data_file: QcowRawFile,
backing_file: Option<Arc<dyn BackingRead>>,
sparse: bool,
/// O_DIRECT alignment requirement (0 = no alignment needed).
alignment: usize,
/// I/O alignment for the AsyncIo trait (at least SECTOR_SIZE).
io_alignment: u64,
cluster_size: u64,
decoder: Arc<dyn Decoder>,
io_uring: IoUring,
eventfd: EventFd,
completion_list: VecDeque<(u64, i32)>,
}
impl QcowAsync {
pub(crate) fn new(
metadata: Arc<QcowMetadata>,
data_file: QcowRawFile,
backing_file: Option<Arc<dyn BackingRead>>,
sparse: bool,
ring_depth: u32,
) -> io::Result<Self> {
let alignment = data_file.file().alignment();
let io_alignment = max(alignment as u64, SECTOR_SIZE);
let io_uring = IoUring::new(ring_depth)?;
let eventfd = EventFd::new(libc::EFD_NONBLOCK)?;
io_uring.submitter().register_eventfd(eventfd.as_raw_fd())?;
Ok(QcowAsync {
cluster_size: metadata.cluster_size(),
decoder: metadata.decoder(),
metadata,
data_file,
backing_file,
sparse,
alignment,
io_alignment,
io_uring,
eventfd,
completion_list: VecDeque::new(),
})
}
fn apply_dealloc_action(&mut self, action: &DeallocAction) {
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);
}
}
}
}
impl AsyncIo for QcowAsync {
fn notifier(&self) -> &EventFd {
&self.eventfd
}
fn read_vectored(
&mut self,
offset: libc::off_t,
iovecs: &[libc::iovec],
user_data: u64,
) -> AsyncIoResult<()> {
let total_len: usize = iovecs.iter().map(|v| v.iov_len).sum();
if let Some(host_offset) = Self::resolve_read(
&self.metadata,
&self.data_file,
&self.backing_file,
offset as u64,
iovecs,
total_len,
self.alignment,
self.cluster_size,
&*self.decoder,
)? {
let fd = self.data_file.as_raw_fd();
let (submitter, mut sq, _) = self.io_uring.split();
// SAFETY: fd is valid and iovecs point to valid guest memory.
unsafe {
sq.push(
&opcode::Readv::new(types::Fd(fd), iovecs.as_ptr(), iovecs.len() as u32)
.offset(host_offset)
.build()
.user_data(user_data),
)
.map_err(|_| {
AsyncIoError::ReadVectored(io::Error::other("Submission queue is full"))
})?;
};
sq.sync();
submitter.submit().map_err(AsyncIoError::ReadVectored)?;
} else {
self.completion_list
.push_back((user_data, total_len as i32));
self.eventfd.write(1).unwrap();
}
Ok(())
}
// TODO Make writes async.
// Writes are synchronous. Async writes require a multi step
// state machine for COW (backing read, cluster allocation, data
// write, L2 commit) with per request buffer lifetime tracking
// and write ordering.
fn write_vectored(
&mut self,
offset: libc::off_t,
iovecs: &[libc::iovec],
user_data: u64,
) -> AsyncIoResult<()> {
Self::cow_write_sync(
offset as u64,
iovecs,
&self.metadata,
&self.data_file,
&self.backing_file,
self.alignment,
self.cluster_size,
)?;
let total_len: usize = iovecs.iter().map(|v| v.iov_len).sum();
self.completion_list
.push_back((user_data, total_len as i32));
self.eventfd.write(1).unwrap();
Ok(())
}
fn fsync(&mut self, user_data: Option<u64>) -> AsyncIoResult<()> {
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)> {
// Drain io_uring completions first, then synthetic ones.
self.io_uring
.completion()
.next()
.map(|entry| (entry.user_data(), entry.result()))
.or_else(|| self.completion_list.pop_front())
}
fn punch_hole(&mut self, offset: u64, length: u64, user_data: u64) -> AsyncIoResult<()> {
let virtual_size = self.metadata.virtual_size();
let cluster_size = self.cluster_size;
let result = self
.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(actions) => {
for action in &actions {
self.apply_dealloc_action(action);
}
self.completion_list.push_back((user_data, 0));
self.eventfd.write(1).unwrap();
Ok(())
}
Err(e) => {
let errno = if let AsyncIoError::PunchHole(ref io_err) = e {
-io_err.raw_os_error().unwrap_or(libc::EIO)
} else {
-libc::EIO
};
self.completion_list.push_back((user_data, errno));
self.eventfd.write(1).unwrap();
Ok(())
}
}
}
fn write_zeroes(&mut self, offset: u64, length: u64, user_data: u64) -> AsyncIoResult<()> {
// For QCOW2, zeroing and hole punching are the same operation.
// Both discard guest data so the range reads back as zero.
self.punch_hole(offset, length, user_data)
}
fn batch_requests_enabled(&self) -> bool {
true
}
fn alignment(&self) -> u64 {
self.io_alignment
}
fn submit_batch_requests(&mut self, batch_request: &[BatchRequest]) -> AsyncIoResult<()> {
let (submitter, mut sq, _) = self.io_uring.split();
let mut needs_submit = false;
let mut sync_completions: Vec<(u64, i32)> = Vec::new();
for req in batch_request {
match req.request_type {
RequestType::In => {
let total_len: usize = req.iovecs.iter().map(|v| v.iov_len).sum();
if let Some(host_offset) = Self::resolve_read(
&self.metadata,
&self.data_file,
&self.backing_file,
req.offset as u64,
&req.iovecs,
total_len,
self.alignment,
self.cluster_size,
&*self.decoder,
)? {
let fd = self.data_file.as_raw_fd();
// SAFETY: fd is valid and iovecs point to valid guest memory.
unsafe {
sq.push(
&opcode::Readv::new(
types::Fd(fd),
req.iovecs.as_ptr(),
req.iovecs.len() as u32,
)
.offset(host_offset)
.build()
.user_data(req.user_data),
)
.map_err(|_| {
AsyncIoError::ReadVectored(io::Error::other(
"Submission queue is full",
))
})?;
}
needs_submit = true;
} else {
sync_completions.push((req.user_data, total_len as i32));
}
}
RequestType::Out => {
let total_len: usize = req.iovecs.iter().map(|v| v.iov_len).sum();
Self::cow_write_sync(
req.offset as u64,
&req.iovecs,
&self.metadata,
&self.data_file,
&self.backing_file,
self.alignment,
self.cluster_size,
)?;
sync_completions.push((req.user_data, total_len as i32));
}
_ => {
unreachable!("Unexpected batch request type: {:?}", req.request_type)
}
}
}
if needs_submit {
sq.sync();
submitter
.submit()
.map_err(AsyncIoError::SubmitBatchRequests)?;
}
if !sync_completions.is_empty() {
for c in sync_completions {
self.completion_list.push_back(c);
}
self.eventfd.write(1).unwrap();
}
Ok(())
}
}
impl QcowAsync {
/// Resolves read mappings for a guest read request.
///
/// Returns `Some(host_offset)` if the entire read falls within a single
/// allocated cluster (fast path). Otherwise handles the read
/// synchronously via `scatter_read_sync` and returns `None`.
#[allow(clippy::too_many_arguments)]
fn resolve_read(
metadata: &QcowMetadata,
data_file: &QcowRawFile,
backing_file: &Option<Arc<dyn BackingRead>>,
address: u64,
iovecs: &[libc::iovec],
total_len: usize,
alignment: usize,
cluster_size: u64,
decoder: &dyn Decoder,
) -> AsyncIoResult<Option<u64>> {
let has_backing = backing_file.is_some();
let mappings = metadata
.map_clusters_for_read(address, total_len, has_backing)
.map_err(AsyncIoError::ReadVectored)?;
// The fast path returns a host offset so the caller can submit a
// single io_uring readv with the original iovecs. This only works
// without O_DIRECT because it requires I/O
// size and file offset to be multiples of the device sector size.
// Guest requests can be smaller (e.g. 512 byte UEFI reads on a
// 4096 byte sector device), so O_DIRECT reads fall through to the
// alignment aware synchronous path instead.
if alignment == 0
&& mappings.len() == 1
&& let ClusterReadMapping::Allocated {
offset: host_offset,
length,
} = &mappings[0]
&& *length as usize == total_len
{
return Ok(Some(*host_offset));
}
Self::scatter_read_sync(
mappings,
iovecs,
data_file,
backing_file,
alignment,
cluster_size,
decoder,
)?;
Ok(None)
}
/// Scatter-read cluster mappings synchronously into iovec buffers.
fn scatter_read_sync(
mappings: Vec<ClusterReadMapping>,
iovecs: &[libc::iovec],
data_file: &QcowRawFile,
backing_file: &Option<Arc<dyn BackingRead>>,
alignment: usize,
cluster_size: u64,
decoder: &dyn Decoder,
) -> AsyncIoResult<()> {
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 len = length as usize;
if alignment > 0 {
let mut abuf =
AlignedBuf::new(len, alignment).map_err(AsyncIoError::ReadVectored)?;
aligned_pread(
data_file.as_raw_fd(),
abuf.as_mut_slice(len),
host_offset,
alignment,
)
.map_err(AsyncIoError::ReadVectored)?;
// SAFETY: iovecs point to valid guest memory buffers.
unsafe { scatter_to_iovecs(iovecs, buf_offset, abuf.as_slice(len)) };
} else {
let mut buf = vec![0u8; len];
pread_exact(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 += len;
}
ClusterReadMapping::Compressed {
host_offset,
compressed_size,
cluster_offset,
length,
} => {
let compressed =
pread_alloc(data_file.as_raw_fd(), host_offset, compressed_size)
.map_err(AsyncIoError::ReadVectored)?;
let decompressed =
decompress_cluster(&compressed, cluster_size as usize, decoder)
.map_err(AsyncIoError::ReadVectored)?;
// SAFETY: iovecs point to valid guest memory buffers.
unsafe {
scatter_to_iovecs(
iovecs,
buf_offset,
&decompressed[cluster_offset..cluster_offset + length],
);
}
buf_offset += length;
}
ClusterReadMapping::Backing {
offset: backing_offset,
length,
} => {
let mut buf = vec![0u8; length as usize];
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;
}
}
}
Ok(())
}
/// Write iovec data cluster-by-cluster with COW from backing file.
fn cow_write_sync(
address: u64,
iovecs: &[libc::iovec],
metadata: &QcowMetadata,
data_file: &QcowRawFile,
backing_file: &Option<Arc<dyn BackingRead>>,
alignment: usize,
cluster_size: u64,
) -> AsyncIoResult<()> {
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 intra_offset = curr_addr & (cluster_size - 1);
let remaining_in_cluster = (cluster_size - intra_offset) as usize;
let count = min(total_len - buf_offset, remaining_in_cluster);
let backing_data = if let Some(backing) = 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 = metadata
.map_cluster_for_write(curr_addr, backing_data)
.map_err(AsyncIoError::WriteVectored)?;
match mapping {
ClusterWriteMapping::Allocated {
offset: host_offset,
} => {
if alignment > 0 {
// O_DIRECT, gather directly into aligned buffer.
let mut abuf = AlignedBuf::new(count, alignment)
.map_err(AsyncIoError::WriteVectored)?;
// SAFETY: iovecs point to valid guest memory buffers
unsafe {
gather_from_iovecs_into(iovecs, buf_offset, abuf.as_mut_slice(count));
}
aligned_pwrite(
data_file.as_raw_fd(),
abuf.as_slice(count),
host_offset,
alignment,
)
.map_err(AsyncIoError::WriteVectored)?;
} else {
// No O_DIRECT, plain buffer is fine.
let mut buf = vec![0u8; count];
// SAFETY: iovecs point to valid guest memory buffers.
unsafe {
gather_from_iovecs_into(iovecs, buf_offset, &mut buf);
}
pwrite_all(data_file.as_raw_fd(), &buf, host_offset)
.map_err(AsyncIoError::WriteVectored)?;
}
}
}
buf_offset += count;
}
Ok(())
}
}
#[cfg(test)]
mod unit_tests {
use std::io::{Seek, SeekFrom, Write};
use std::sync::Arc;
use std::thread;
use vmm_sys_util::tempfile::TempFile;
use super::*;
use crate::disk_file::AsyncDiskFile;
use crate::qcow::{QcowFile, RawFile};
use crate::qcow_common::unit_tests::compress_allocated_clusters;
use crate::qcow_disk::QcowDisk;
use crate::{BatchRequest, RequestType, SECTOR_SIZE};
fn create_disk_with_data(
file_size: u64,
data: &[u8],
offset: u64,
sparse: bool,
) -> (TempFile, QcowDisk) {
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 = QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
false,
false,
sparse,
true,
)
.unwrap();
(temp_file, disk)
}
fn wait_for_completion(async_io: &mut dyn AsyncIo) -> (u64, i32) {
loop {
if let Some(c) = async_io.next_completed_request() {
return c;
}
// Block until the eventfd is signaled (io_uring or synthetic).
let fd = async_io.notifier().as_raw_fd();
let mut val = 0u64;
// SAFETY: reading 8 bytes from a valid eventfd.
unsafe {
libc::read(fd, (&raw mut val).cast(), 8);
}
}
}
fn async_write(disk: &QcowDisk, offset: u64, data: &[u8]) {
let mut async_io = disk.create_async_io(1).unwrap();
let iovec = libc::iovec {
iov_base: data.as_ptr().cast::<libc::c_void>().cast_mut(),
iov_len: data.len(),
};
async_io
.write_vectored(offset as libc::off_t, &[iovec], 2)
.unwrap();
let (user_data, result) = wait_for_completion(async_io.as_mut());
assert_eq!(user_data, 2);
assert_eq!(
result as usize,
data.len(),
"write should return requested length"
);
}
fn async_read(disk: &QcowDisk, offset: u64, len: usize) -> Vec<u8> {
let mut async_io = disk.create_async_io(1).unwrap();
let mut buf = vec![0xFFu8; len];
let iovec = libc::iovec {
iov_base: buf.as_mut_ptr().cast(),
iov_len: buf.len(),
};
async_io
.read_vectored(offset as libc::off_t, &[iovec], 1)
.unwrap();
let (user_data, result) = wait_for_completion(async_io.as_mut());
assert_eq!(user_data, 1);
assert_eq!(result as usize, len, "read should return requested length");
buf
}
#[test]
fn test_qcow_async_punch_hole_completion() {
let data = vec![0xDD; 128 * 1024];
let offset = 0u64;
let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true);
let mut async_io = disk.create_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);
let read_buf = async_read(&disk, offset, data.len());
assert!(
read_buf.iter().all(|&b| b == 0),
"Punched hole should read as zeros"
);
}
#[test]
fn test_qcow_async_write_zeroes_completion() {
let data = vec![0xAA; 128 * 1024];
let offset = 0u64;
let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true);
let mut async_io = disk.create_async_io(1).unwrap();
async_io
.write_zeroes(offset, data.len() as u64, 200)
.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);
let read_buf = async_read(&disk, offset, data.len());
assert!(
read_buf.iter().all(|&b| b == 0),
"Write zeroes region should read as zeros"
);
}
#[test]
fn test_qcow_async_write_read_roundtrip() {
let file_size = 100 * 1024 * 1024;
let temp_file = TempFile::new().unwrap();
{
let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false);
QcowFile::new(raw_file, 3, file_size, true).unwrap();
}
let disk = QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
false,
false,
true,
true,
)
.unwrap();
let pattern: Vec<u8> = (0..128 * 1024).map(|i| (i % 251) as u8).collect();
let offset = 64 * 1024;
async_write(&disk, offset, &pattern);
let read_buf = async_read(&disk, offset, pattern.len());
assert_eq!(read_buf, pattern, "read should match written data");
}
#[test]
fn test_qcow_async_read_spanning_cluster_boundary() {
let cluster_size: u64 = 65536;
let file_size = 100 * 1024 * 1024;
// Write distinct patterns into two adjacent clusters.
let pattern_a = vec![0xAA; cluster_size as usize];
let pattern_b = vec![0xBB; cluster_size as usize];
let (_temp, disk) = create_disk_with_data(file_size, &pattern_a, 0, true);
async_write(&disk, cluster_size, &pattern_b);
// Read across the boundary: last 4K of cluster 0 + first 4K of cluster 1.
let read_offset = cluster_size - 4096;
let read_len = 8192;
let buf = async_read(&disk, read_offset, read_len);
assert!(
buf[..4096].iter().all(|&b| b == 0xAA),
"first half should come from cluster 0"
);
assert!(
buf[4096..].iter().all(|&b| b == 0xBB),
"second half should come from cluster 1"
);
}
#[test]
fn test_qcow_async_batch_mixed_requests() {
let file_size = 100 * 1024 * 1024;
let temp_file = TempFile::new().unwrap();
{
let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false);
QcowFile::new(raw_file, 3, file_size, true).unwrap();
}
let disk = QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
false,
false,
true,
true,
)
.unwrap();
let mut async_io = disk.create_async_io(8).unwrap();
// Prepare write data for two regions.
let write_a = vec![0xAA; 4096];
let write_b = vec![0xBB; 4096];
let offset_a: u64 = 0;
let offset_b: u64 = 65536;
let iov_a = libc::iovec {
iov_base: write_a.as_ptr().cast::<libc::c_void>().cast_mut(),
iov_len: write_a.len(),
};
let iov_b = libc::iovec {
iov_base: write_b.as_ptr().cast::<libc::c_void>().cast_mut(),
iov_len: write_b.len(),
};
let batch = vec![
BatchRequest {
offset: offset_a as libc::off_t,
iovecs: smallvec::smallvec![iov_a],
user_data: 10,
request_type: RequestType::Out,
},
BatchRequest {
offset: offset_b as libc::off_t,
iovecs: smallvec::smallvec![iov_b],
user_data: 20,
request_type: RequestType::Out,
},
];
async_io.submit_batch_requests(&batch).unwrap();
let mut completions = [
wait_for_completion(async_io.as_mut()),
wait_for_completion(async_io.as_mut()),
];
completions.sort_by_key(|c| c.0);
assert_eq!(completions[0], (10, 4096));
assert_eq!(completions[1], (20, 4096));
drop(async_io);
// Batch read both regions back.
let mut read_a = vec![0u8; 4096];
let mut read_b = vec![0u8; 4096];
let riov_a = libc::iovec {
iov_base: read_a.as_mut_ptr().cast(),
iov_len: read_a.len(),
};
let riov_b = libc::iovec {
iov_base: read_b.as_mut_ptr().cast(),
iov_len: read_b.len(),
};
let mut async_io = disk.create_async_io(8).unwrap();
let read_batch = vec![
BatchRequest {
offset: offset_a as libc::off_t,
iovecs: smallvec::smallvec![riov_a],
user_data: 30,
request_type: RequestType::In,
},
BatchRequest {
offset: offset_b as libc::off_t,
iovecs: smallvec::smallvec![riov_b],
user_data: 40,
request_type: RequestType::In,
},
];
async_io.submit_batch_requests(&read_batch).unwrap();
let mut completions = [
wait_for_completion(async_io.as_mut()),
wait_for_completion(async_io.as_mut()),
];
completions.sort_by_key(|c| c.0);
assert_eq!(completions[0], (30, 4096));
assert_eq!(completions[1], (40, 4096));
assert_eq!(read_a, write_a, "batch read A should match written data");
assert_eq!(read_b, write_b, "batch read B should match written data");
}
#[test]
fn test_qcow_async_read_unallocated() {
let file_size = 100 * 1024 * 1024;
let temp_file = TempFile::new().unwrap();
{
let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false);
QcowFile::new(raw_file, 3, file_size, true).unwrap();
}
let disk = QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
false,
false,
true,
true,
)
.unwrap();
let buf = async_read(&disk, 0, 128 * 1024);
assert!(
buf.iter().all(|&b| b == 0),
"unallocated region should read as zeroes"
);
}
#[test]
fn test_qcow_async_sub_cluster_write() {
let cluster_size = 65536usize;
let file_size = 100 * 1024 * 1024;
let temp_file = TempFile::new().unwrap();
{
let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false);
QcowFile::new(raw_file, 3, file_size, true).unwrap();
}
let disk = QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
false,
false,
true,
true,
)
.unwrap();
// Write 4K into the middle of a cluster.
let write_offset = 4096u64;
let write_len = 4096;
let pattern = vec![0xCC; write_len];
async_write(&disk, write_offset, &pattern);
// Read the entire cluster back.
let buf = async_read(&disk, 0, cluster_size);
assert!(
buf[..write_offset as usize].iter().all(|&b| b == 0),
"bytes before the write should be zero"
);
assert_eq!(
&buf[write_offset as usize..write_offset as usize + write_len],
&pattern[..],
"written region should match"
);
assert!(
buf[write_offset as usize + write_len..]
.iter()
.all(|&b| b == 0),
"bytes after the write should be zero"
);
}
#[test]
fn test_qcow_async_write_after_punch_hole() {
let data = vec![0xAA; 64 * 1024];
let offset = 0u64;
let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true);
let buf = async_read(&disk, offset, data.len());
assert!(buf.iter().all(|&b| b == 0xAA));
let mut async_io = disk.create_async_io(1).unwrap();
async_io.punch_hole(offset, data.len() as u64, 10).unwrap();
let (_, result) = wait_for_completion(async_io.as_mut());
assert_eq!(result, 0);
drop(async_io);
let buf = async_read(&disk, offset, data.len());
assert!(
buf.iter().all(|&b| b == 0),
"should be zero after punch hole"
);
let new_data = vec![0xBB; 64 * 1024];
async_write(&disk, offset, &new_data);
let buf = async_read(&disk, offset, new_data.len());
assert_eq!(buf, new_data, "should read new data after rewrite");
}
#[test]
fn test_qcow_async_large_sequential_io() {
let cluster_size = 64 * 1024;
let num_clusters = 8;
let total_len = cluster_size * num_clusters;
let offset = 0u64;
let mut data = vec![0u8; total_len];
for (i, chunk) in data.chunks_mut(cluster_size).enumerate() {
chunk.fill((i + 1) as u8);
}
let (_temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, offset, true);
let buf = async_read(&disk, offset, total_len);
assert_eq!(buf.len(), total_len);
for (i, chunk) in buf.chunks(cluster_size).enumerate() {
assert!(
chunk.iter().all(|&b| b == (i + 1) as u8),
"cluster {i} mismatch"
);
}
}
#[test]
fn test_qcow_async_alignment_without_direct_io() {
let file_size = 100 * 1024 * 1024;
let temp_file = TempFile::new().unwrap();
{
let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false);
QcowFile::new(raw_file, 3, file_size, true).unwrap();
}
let disk = QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
false,
false,
true,
true,
)
.unwrap();
let async_io = disk.create_async_io(1).unwrap();
assert_eq!(async_io.alignment(), SECTOR_SIZE);
}
/// Returns None if O_DIRECT is not supported (e.g. tmpfs).
fn try_create_direct_io_disk(temp_file: &TempFile, file_size: u64) -> Option<QcowDisk> {
{
let raw_file = RawFile::new(temp_file.as_file().try_clone().unwrap(), false);
QcowFile::new(raw_file, 3, file_size, true).unwrap();
}
QcowDisk::new(
temp_file.as_file().try_clone().unwrap(),
true,
false,
true,
true,
)
.ok()
}
#[test]
fn test_qcow_async_alignment_with_direct_io() {
let temp_file = TempFile::new().unwrap();
let disk = match try_create_direct_io_disk(&temp_file, 100 * 1024 * 1024) {
Some(d) => d,
None => {
eprintln!("skipping: O_DIRECT not supported on this filesystem");
return;
}
};
let async_io = disk.create_async_io(1).unwrap();
assert!(async_io.alignment() >= SECTOR_SIZE);
}
#[test]
fn test_qcow_async_sub_sector_read_with_direct_io() {
let temp_file = TempFile::new().unwrap();
let disk = match try_create_direct_io_disk(&temp_file, 100 * 1024 * 1024) {
Some(d) => d,
None => {
eprintln!("skipping: O_DIRECT not supported on this filesystem");
return;
}
};
let pattern = vec![0xAB; 65536];
async_write(&disk, 0, &pattern);
let buf = async_read(&disk, 0, 512);
assert!(
buf.iter().all(|&b| b == 0xAB),
"sub-sector O_DIRECT read should return written data"
);
}
#[test]
fn test_qcow_async_direct_io_write_read_roundtrip() {
let temp_file = TempFile::new().unwrap();
let disk = match try_create_direct_io_disk(&temp_file, 100 * 1024 * 1024) {
Some(d) => d,
None => {
eprintln!("skipping: O_DIRECT not supported on this filesystem");
return;
}
};
let pattern: Vec<u8> = (0..128 * 1024).map(|i| (i % 251) as u8).collect();
async_write(&disk, 0, &pattern);
let buf = async_read(&disk, 0, pattern.len());
assert_eq!(buf, pattern, "O_DIRECT roundtrip should match");
}
#[test]
fn test_compressed_read_multi_queue() {
let cluster_size = 65536usize;
let data: Vec<u8> = (0..=255).cycle().take(cluster_size).collect();
let (temp, disk) = create_disk_with_data(100 * 1024 * 1024, &data, 0, false);
drop(disk);
compress_allocated_clusters(&mut temp.as_file().try_clone().unwrap());
let disk = Arc::new(
QcowDisk::new(
temp.as_file().try_clone().unwrap(),
false,
false,
false,
true,
)
.unwrap(),
);
let handles: Vec<_> = (0..4)
.map(|_| {
let disk = Arc::clone(&disk);
let expected = data.clone();
thread::spawn(move || {
let mut async_io = disk.create_async_io(1).unwrap();
let mut buf = vec![0xFFu8; cluster_size];
let iovec = libc::iovec {
iov_base: buf.as_mut_ptr().cast(),
iov_len: buf.len(),
};
async_io.read_vectored(0, &[iovec], 1).unwrap();
let (_, result) = wait_for_completion(async_io.as_mut());
assert_eq!(result as usize, cluster_size);
assert_eq!(buf, expected);
})
})
.collect();
for h in handles {
h.join().unwrap();
}
}
}