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); }) })