From afdc87e97172cc5582334bbdecf91f72ca9a8138 Mon Sep 17 00:00:00 2001 From: Dylan Reid Date: Tue, 19 May 2026 14:37:33 -0700 Subject: [PATCH] block: Add owned qcow io_uring I/O path Similar to the previous commits, use UringDataIo for qcow async. Again, the legacy interfaces are kept(at the expense of some temporary code). The temporary code is unsound, like the existing code, but will be removed soon. Signed-off-by: Dylan Reid --- block/src/qcow_async.rs | 448 ++++++++++++++++++++++++---------------- 1 file changed, 269 insertions(+), 179 deletions(-) diff --git a/block/src/qcow_async.rs b/block/src/qcow_async.rs index 0019c80e6..680aa8f77 100644 --- a/block/src/qcow_async.rs +++ b/block/src/qcow_async.rs @@ -2,21 +2,23 @@ // // Copyright 2026 The Cloud Hypervisor Authors. All rights reserved. // +// Copyright (c) Meta Platforms, Inc. and affiliates. +// // 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::async_io::{ + AsyncIo, AsyncIoCompletion, AsyncIoError, AsyncIoOperation, AsyncIoResult, UringDataIo, +}; use crate::qcow::decoder::Decoder; use crate::qcow::metadata::{ BackingRead, ClusterReadMapping, ClusterWriteMapping, DeallocAction, QcowMetadata, @@ -48,9 +50,7 @@ pub struct QcowAsync { io_alignment: u64, cluster_size: u64, decoder: Arc, - io_uring: IoUring, - eventfd: EventFd, - completion_list: VecDeque<(u64, i32)>, + data_io: UringDataIo, } impl QcowAsync { @@ -63,9 +63,6 @@ impl QcowAsync { ) -> io::Result { 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(), @@ -76,9 +73,7 @@ impl QcowAsync { sparse, alignment, io_alignment, - io_uring, - eventfd, - completion_list: VecDeque::new(), + data_io: UringDataIo::new(ring_depth)?, }) } @@ -101,11 +96,83 @@ impl QcowAsync { } } } + + fn async_error_result(error: &AsyncIoError) -> i32 { + let io_error = match error { + AsyncIoError::ReadVectored(e) + | AsyncIoError::WriteVectored(e) + | AsyncIoError::SubmitBatchRequests(e) + | AsyncIoError::Fsync(e) + | AsyncIoError::PunchHole(e) + | AsyncIoError::WriteZeroes(e) => e, + }; + -io_error.raw_os_error().unwrap_or(libc::EIO) + } + + fn inject_operation_completion(&mut self, op: AsyncIoOperation, result: i32) { + self.data_io + .inject_completion(AsyncIoCompletion::from_operation(op, result)); + } + + fn prepare_read_operation( + &mut self, + mut op: AsyncIoOperation, + ) -> Result, Box<(AsyncIoOperation, AsyncIoError)>> { + let total_len = op.total_len(); + let host_offset = match Self::resolve_read( + &self.metadata, + &self.data_file, + &self.backing_file, + op.offset() as u64, + op.iovecs(), + total_len, + self.alignment, + self.cluster_size, + &*self.decoder, + ) { + Ok(host_offset) => host_offset, + Err(e) => return Err(Box::new((op, e))), + }; + + if let Some(host_offset) = host_offset { + op.set_offset(host_offset as libc::off_t); + Ok(Some(op)) + } else { + self.inject_operation_completion(op, total_len as i32); + Ok(None) + } + } + + fn complete_write_operation_sync( + &mut self, + op: AsyncIoOperation, + ) -> Result<(), Box<(AsyncIoOperation, AsyncIoError)>> { + // 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. + let total_len = op.total_len(); + if let Err(e) = Self::cow_write_sync( + op.offset() as u64, + op.iovecs(), + &self.metadata, + &self.data_file, + &self.backing_file, + self.alignment, + self.cluster_size, + ) { + return Err(Box::new((op, e))); + } + + self.inject_operation_completion(op, total_len as i32); + Ok(()) + } } impl AsyncIo for QcowAsync { fn notifier(&self) -> &EventFd { - &self.eventfd + self.data_io.notifier() } fn read_vectored( @@ -127,37 +194,28 @@ impl AsyncIo for QcowAsync { 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. + // SAFETY: this legacy trait method's caller must keep the + // borrowed iovecs and writable buffers valid until completion. unsafe { - sq.push( - &opcode::Readv::new(types::Fd(fd), iovecs.as_ptr(), iovecs.len() as u32) - .offset(host_offset) - .build() - .user_data(user_data), + self.data_io.submit_borrowed_operation( + self.data_file.as_raw_fd(), + host_offset as libc::off_t, + true, + iovecs, + user_data, ) - .map_err(|_| { - AsyncIoError::ReadVectored(io::Error::other("Submission queue is full")) - })?; - }; - - sq.sync(); - submitter.submit().map_err(AsyncIoError::ReadVectored)?; + } + .map_err(AsyncIoError::ReadVectored)?; } else { - self.completion_list - .push_back((user_data, total_len as i32)); - self.eventfd.write(1).unwrap(); + self.data_io.inject_completion(AsyncIoCompletion::new( + user_data, + total_len as i32, + None, + )); } 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, @@ -175,28 +233,39 @@ impl AsyncIo for QcowAsync { )?; 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(); + self.data_io + .inject_completion(AsyncIoCompletion::new(user_data, total_len as i32, None)); Ok(()) } + fn submit_data_operation(&mut self, op: AsyncIoOperation) -> AsyncIoResult<()> { + if op.is_read() { + match self.prepare_read_operation(op) { + Ok(Some(op)) => { + self.data_io + .submit_operation(self.data_file.as_raw_fd(), op) + .map_err(AsyncIoError::ReadVectored)?; + } + Ok(None) => {} + Err(e) => return Err(e.1), + } + Ok(()) + } else { + self.complete_write_operation_sync(op).map_err(|e| e.1) + } + } + fn fsync(&mut self, user_data: Option) -> 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(); + self.data_io + .inject_completion(AsyncIoCompletion::new(user_data, 0, None)); } 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 next_completion(&mut self) -> Option { + self.data_io.next_completion() } fn punch_hole(&mut self, offset: u64, length: u64, user_data: u64) -> AsyncIoResult<()> { @@ -216,8 +285,8 @@ impl AsyncIo for QcowAsync { for action in &actions { self.apply_dealloc_action(action); } - self.completion_list.push_back((user_data, 0)); - self.eventfd.write(1).unwrap(); + self.data_io + .inject_completion(AsyncIoCompletion::new(user_data, 0, None)); Ok(()) } Err(e) => { @@ -226,8 +295,8 @@ impl AsyncIo for QcowAsync { } else { -libc::EIO }; - self.completion_list.push_back((user_data, errno)); - self.eventfd.write(1).unwrap(); + self.data_io + .inject_completion(AsyncIoCompletion::new(user_data, errno, None)); Ok(()) } } @@ -250,8 +319,8 @@ impl AsyncIo for QcowAsync { for action in &actions { self.apply_dealloc_action(action); } - self.completion_list.push_back((user_data, 0)); - self.eventfd.write(1).unwrap(); + self.data_io + .inject_completion(AsyncIoCompletion::new(user_data, 0, None)); Ok(()) } Err(e) => { @@ -260,8 +329,8 @@ impl AsyncIo for QcowAsync { } else { -libc::EIO }; - self.completion_list.push_back((user_data, errno)); - self.eventfd.write(1).unwrap(); + self.data_io + .inject_completion(AsyncIoCompletion::new(user_data, errno, None)); Ok(()) } } @@ -276,9 +345,7 @@ impl AsyncIo for QcowAsync { } 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(); + let mut async_reads = Vec::new(); for req in batch_request { match req.request_type { @@ -296,28 +363,18 @@ impl AsyncIo for QcowAsync { 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; + async_reads.push(( + host_offset as libc::off_t, + true, + req.iovecs.as_slice(), + req.user_data, + )); } else { - sync_completions.push((req.user_data, total_len as i32)); + self.data_io.inject_completion(AsyncIoCompletion::new( + req.user_data, + total_len as i32, + None, + )); } } RequestType::Out => { @@ -331,28 +388,63 @@ impl AsyncIo for QcowAsync { self.alignment, self.cluster_size, )?; - sync_completions.push((req.user_data, total_len as i32)); - } - _ => { - unreachable!("Unexpected batch request type: {:?}", req.request_type) + self.data_io.inject_completion(AsyncIoCompletion::new( + req.user_data, + total_len as i32, + None, + )); } + _ => unreachable!("Unexpected batch request type: {:?}", req.request_type), } } - if needs_submit { - sq.sync(); - submitter - .submit() + if !async_reads.is_empty() { + // SAFETY: this legacy trait method's caller must keep every + // borrowed iovec array and buffer valid until its completion. + unsafe { + self.data_io + .submit_borrowed_batch(self.data_file.as_raw_fd(), &async_reads) + } + .map_err(AsyncIoError::SubmitBatchRequests)?; + } + + Ok(()) + } + + fn submit_batch_operations( + &mut self, + batch_request: Vec, + ) -> AsyncIoResult<()> { + let mut async_reads = Vec::new(); + + for op in batch_request { + if op.is_read() { + match self.prepare_read_operation(op) { + Ok(Some(op)) => async_reads.push(op), + Ok(None) => {} + Err(boxed) => { + let (op, e) = *boxed; + // The operation was not submitted to the kernel. Accept + // it at the qcow layer and surface the failure through + // the common completion path so batch acceptance remains + // all-or-none for the virtqueue. + let result = Self::async_error_result(&e); + self.inject_operation_completion(op, result); + } + } + } else if let Err(boxed) = self.complete_write_operation_sync(op) { + let (op, e) = *boxed; + let result = Self::async_error_result(&e); + self.inject_operation_completion(op, result); + } + } + + if !async_reads.is_empty() { + self.data_io + .submit_batch(self.data_file.as_raw_fd(), async_reads) .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(()) } } @@ -580,11 +672,12 @@ mod unit_tests { use vmm_sys_util::tempfile::TempFile; use super::*; + use crate::SECTOR_SIZE; + use crate::async_io::{AsyncIoCompletion, AsyncIoOperation, OwnedIoBuffer}; use crate::disk_file::AsyncDiskFile; use crate::qcow::{BackingFileConfig, ImageType, 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, @@ -643,9 +736,9 @@ mod unit_tests { (backing_temp, overlay_temp, disk) } - fn wait_for_completion(async_io: &mut dyn AsyncIo) -> (u64, i32) { + fn wait_for_completion(async_io: &mut dyn AsyncIo) -> AsyncIoCompletion { loop { - if let Some(c) = async_io.next_completed_request() { + if let Some(c) = async_io.next_completion() { return c; } // Block until the eventfd is signaled (io_uring or synthetic). @@ -658,16 +751,21 @@ mod unit_tests { } } + fn completion_tuple(completion: &AsyncIoCompletion) -> (u64, i32) { + (completion.user_data, completion.result) + } + 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::().cast_mut(), - iov_len: data.len(), - }; async_io - .write_vectored(offset as libc::off_t, &[iovec], 2) + .write_from_vec( + offset as libc::off_t, + OwnedIoBuffer::from_vec(data.to_vec()), + 2, + ) .unwrap(); - let (user_data, result) = wait_for_completion(async_io.as_mut()); + let completion = wait_for_completion(async_io.as_mut()); + let (user_data, result) = completion_tuple(&completion); assert_eq!(user_data, 2); assert_eq!( result as usize, @@ -678,18 +776,21 @@ mod unit_tests { fn async_read(disk: &QcowDisk, offset: u64, len: usize) -> Vec { 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) + .read_to_vec( + offset as libc::off_t, + OwnedIoBuffer::from_vec(vec![0xFF; len]), + 1, + ) .unwrap(); - let (user_data, result) = wait_for_completion(async_io.as_mut()); + let mut completion = wait_for_completion(async_io.as_mut()); + let (user_data, result) = completion_tuple(&completion); assert_eq!(user_data, 1); assert_eq!(result as usize, len, "read should return requested length"); - buf + match completion.buffer.take() { + Some(buffer) => buffer.as_slice().to_vec(), + other => panic!("unexpected read completion: {other:?}"), + } } #[test] @@ -700,7 +801,8 @@ mod unit_tests { 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(); + let completion = async_io.next_completion().unwrap(); + let (user_data, result) = completion_tuple(&completion); assert_eq!(user_data, 100); assert_eq!(result, 0, "punch_hole should succeed"); drop(async_io); @@ -722,7 +824,8 @@ mod unit_tests { async_io .write_zeroes(offset, data.len() as u64, 200) .unwrap(); - let (user_data, result) = async_io.next_completed_request().unwrap(); + let completion = async_io.next_completion().unwrap(); + let (user_data, result) = completion_tuple(&completion); assert_eq!(user_data, 200); assert_eq!(result, 0, "write_zeroes should succeed"); drop(async_io); @@ -744,7 +847,8 @@ mod unit_tests { let mut async_io = disk.create_async_io(1).unwrap(); async_io.write_zeroes(offset, cluster_size, 201).unwrap(); - let (user_data, result) = wait_for_completion(async_io.as_mut()); + let completion = wait_for_completion(async_io.as_mut()); + let (user_data, result) = completion_tuple(&completion); assert_eq!(user_data, 201); assert_eq!(result, 0, "write_zeroes should succeed"); drop(async_io); @@ -832,35 +936,24 @@ mod unit_tests { let offset_a: u64 = 0; let offset_b: u64 = 65536; - let iov_a = libc::iovec { - iov_base: write_a.as_ptr().cast::().cast_mut(), - iov_len: write_a.len(), - }; - let iov_b = libc::iovec { - iov_base: write_b.as_ptr().cast::().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, - }, + AsyncIoOperation::write_from_vec( + offset_a as libc::off_t, + OwnedIoBuffer::from_vec(write_a.clone()), + 10, + ), + AsyncIoOperation::write_from_vec( + offset_b as libc::off_t, + OwnedIoBuffer::from_vec(write_b.clone()), + 20, + ), ]; - async_io.submit_batch_requests(&batch).unwrap(); + async_io.submit_batch_operations(batch).unwrap(); let mut completions = [ - wait_for_completion(async_io.as_mut()), - wait_for_completion(async_io.as_mut()), + completion_tuple(&wait_for_completion(async_io.as_mut())), + completion_tuple(&wait_for_completion(async_io.as_mut())), ]; completions.sort_by_key(|c| c.0); assert_eq!(completions[0], (10, 4096)); @@ -868,43 +961,38 @@ mod unit_tests { 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, - }, + AsyncIoOperation::read_to_vec( + offset_a as libc::off_t, + OwnedIoBuffer::from_vec(vec![0; 4096]), + 30, + ), + AsyncIoOperation::read_to_vec( + offset_b as libc::off_t, + OwnedIoBuffer::from_vec(vec![0; 4096]), + 40, + ), ]; - async_io.submit_batch_requests(&read_batch).unwrap(); + async_io.submit_batch_operations(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)); + let mut completion_a = wait_for_completion(async_io.as_mut()); + let mut completion_b = wait_for_completion(async_io.as_mut()); + if completion_a.user_data > completion_b.user_data { + std::mem::swap(&mut completion_a, &mut completion_b); + } + assert_eq!(completion_tuple(&completion_a), (30, 4096)); + assert_eq!(completion_tuple(&completion_b), (40, 4096)); + let read_a = match completion_a.buffer.take() { + Some(buffer) => buffer.as_slice().to_vec(), + other => panic!("unexpected read completion A: {other:?}"), + }; + let read_b = match completion_b.buffer.take() { + Some(buffer) => buffer.as_slice().to_vec(), + other => panic!("unexpected read completion B: {other:?}"), + }; 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"); } @@ -988,7 +1076,7 @@ mod unit_tests { 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()); + let result = wait_for_completion(async_io.as_mut()).result; assert_eq!(result, 0); drop(async_io); @@ -1144,14 +1232,16 @@ mod unit_tests { 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()); + async_io + .read_to_vec(0, OwnedIoBuffer::from_vec(vec![0xFF; cluster_size]), 1) + .unwrap(); + let mut completion = wait_for_completion(async_io.as_mut()); + let result = completion.result; assert_eq!(result as usize, cluster_size); + let buf = match completion.buffer.take() { + Some(buffer) => buffer.as_slice().to_vec(), + other => panic!("unexpected read completion: {other:?}"), + }; assert_eq!(buf, expected); }) })