From 914303980584906a3bb3a7c3f1f62b493445aca1 Mon Sep 17 00:00:00 2001 From: Dylan Reid Date: Tue, 19 May 2026 14:37:17 -0700 Subject: [PATCH] block: Add owned raw Linux AIO path Start using AioDataIo from RawFileAsyncAio. This adds a temporary submit_borrowed_operation to enable preserving the unsafe iovec api until we can remove it in the forthcoming commits. Signed-off-by: Dylan Reid --- block/src/async_io/aio_data_io.rs | 88 +++++++++++++++++------- block/src/raw_async_aio.rs | 108 ++++++++++-------------------- 2 files changed, 100 insertions(+), 96 deletions(-) diff --git a/block/src/async_io/aio_data_io.rs b/block/src/async_io/aio_data_io.rs index 410bce56c..d0a97f099 100644 --- a/block/src/async_io/aio_data_io.rs +++ b/block/src/async_io/aio_data_io.rs @@ -6,7 +6,7 @@ // // SPDX-License-Identifier: Apache-2.0 AND BSD-3-Clause -use std::collections::{HashMap, HashSet, VecDeque}; +use std::collections::{HashMap, VecDeque}; use std::io; use std::os::fd::{AsRawFd, RawFd}; @@ -18,24 +18,18 @@ use super::common::{duplicate_user_data_error, errno_result, validate_batch}; use super::{AsyncIoCompletion, AsyncIoOperation}; /// Retained Linux AIO queue for owned async data I/O operations. -/// -/// Submitted operations are kept in `pending` until their completion event is -/// consumed so the iovec pointers handed to the kernel remain valid for the -/// full lifetime of the operation. pub struct AioDataIo { - // Keep this before `pending`: Rust drops fields in declaration order, so + // Keep this before `in_flight`: Rust drops fields in declaration order, so // dropping the context destroys kernel AIO state before retained // operations release the buffers referenced by their iovecs. ctx: aio::IoContext, // The `EventFd` for completion signals. eventfd: EventFd, - // `in_flight` tracks user_data values accepted by the kernel submission - // path, including metadata operations that do not retain a buffer. - in_flight: HashSet, - // `pending` owns read/write operations until their events are consumed so - // their iovecs and backing buffers remain valid while the kernel may use - // them. - pending: HashMap, + // `in_flight` tracks every user_data value accepted by the kernel. Owned + // data operations store `Some(op)` so their iovecs and backing buffers + // remain valid until completion; metadata and legacy borrowed operations + // store `None`. + in_flight: HashMap>, // `completions` holds locally produced completions and kernel events that // have been fetched but not yet returned to the caller. completions: VecDeque, @@ -47,8 +41,7 @@ impl AioDataIo { Ok(Self { ctx: aio::IoContext::new(queue_depth)?, eventfd: EventFd::new(libc::EFD_NONBLOCK)?, - in_flight: HashSet::new(), - pending: HashMap::new(), + in_flight: HashMap::new(), completions: VecDeque::new(), }) } @@ -73,7 +66,7 @@ impl AioDataIo { /// can observe every accepted request through the normal completion path. pub fn submit_operation(&mut self, fd: RawFd, op: AsyncIoOperation) -> io::Result<()> { validate_batch( - |user_data| self.in_flight.contains(&user_data), + |user_data| self.in_flight.contains_key(&user_data), std::slice::from_ref(&op), )?; @@ -95,8 +88,7 @@ impl AioDataIo { aio_resfd: self.eventfd.as_raw_fd() as u32, ..Default::default() }; - self.in_flight.insert(user_data); - self.pending.insert(user_data, op); + self.in_flight.insert(user_data, Some(op)); let result = match Self::submit_iocbs(&self.ctx, &[&mut iocb]) { Ok(1) => return Ok(()), @@ -105,17 +97,65 @@ impl AioDataIo { }; let buffer = self - .pending + .in_flight .remove(&user_data) + .flatten() .and_then(AsyncIoOperation::into_completion_buffer); - self.in_flight.remove(&user_data); self.inject_completion(AsyncIoCompletion::new(user_data, result, buffer)); Ok(()) } + /// Submits one borrowed read or write operation to the queue. + /// + /// This is only for the legacy `AsyncIo` interface. The caller must keep + /// the iovec array and the buffers it references alive until the matching + /// completion is consumed. + pub fn submit_borrowed_operation( + &mut self, + fd: RawFd, + offset: libc::off_t, + is_read: bool, + iovecs: &[libc::iovec], + user_data: u64, + ) -> io::Result<()> { + if self.in_flight.contains_key(&user_data) { + return Err(duplicate_user_data_error(user_data)); + } + self.in_flight.insert(user_data, None); + + let opcode = if is_read { + aio::IOCB_CMD_PREADV + } else { + aio::IOCB_CMD_PWRITEV + }; + let mut iocb = aio::IoControlBlock { + aio_fildes: fd.as_raw_fd() as u32, + aio_lio_opcode: opcode as u16, + aio_buf: iovecs.as_ptr() as u64, + aio_nbytes: iovecs.len() as u64, + aio_offset: offset, + aio_data: user_data, + aio_flags: aio::IOCB_FLAG_RESFD, + aio_resfd: self.eventfd.as_raw_fd() as u32, + ..Default::default() + }; + + match Self::submit_iocbs(&self.ctx, &[&mut iocb]) { + Ok(1) => Ok(()), + Ok(_) => { + self.in_flight.remove(&user_data); + Err(io::Error::from_raw_os_error(libc::EAGAIN)) + } + Err(e) => { + self.in_flight.remove(&user_data); + Err(e) + } + } + } + /// Submits an fsync operation carrying `user_data`. pub fn submit_fsync(&mut self, fd: RawFd, user_data: u64) -> io::Result<()> { - if self.in_flight.contains(&user_data) { + if self.in_flight.contains_key(&user_data) { return Err(duplicate_user_data_error(user_data)); } @@ -127,7 +167,7 @@ impl AioDataIo { aio_resfd: self.eventfd.as_raw_fd() as u32, ..Default::default() }; - self.in_flight.insert(user_data); + self.in_flight.insert(user_data, None); let result = match Self::submit_iocbs(&self.ctx, &[&mut iocb]) { Ok(1) => return Ok(()), Ok(_) => -libc::EAGAIN, @@ -163,12 +203,12 @@ impl AioDataIo { } }; for event in &events[..rc] { - self.in_flight.remove(&event.data); self.completions.push_back(AsyncIoCompletion::new( event.data, event.res as i32, - self.pending + self.in_flight .remove(&event.data) + .flatten() .and_then(AsyncIoOperation::into_completion_buffer), )); } diff --git a/block/src/raw_async_aio.rs b/block/src/raw_async_aio.rs index 22e0896bc..e910bf675 100644 --- a/block/src/raw_async_aio.rs +++ b/block/src/raw_async_aio.rs @@ -1,44 +1,40 @@ // Copyright © 2023 Intel Corporation // +// Copyright (c) Meta Platforms, Inc. and affiliates. +// // SPDX-License-Identifier: Apache-2.0 AND BSD-3-Clause // // Copyright © 2023 Crusoe Energy Systems LLC // -use std::collections::VecDeque; -use std::os::unix::io::{AsRawFd, RawFd}; +use std::os::unix::io::RawFd; -use vmm_sys_util::aio; use vmm_sys_util::eventfd::EventFd; -use crate::async_io::{AsyncIo, AsyncIoError, AsyncIoResult}; +use crate::async_io::{ + AioDataIo, AsyncIo, AsyncIoCompletion, AsyncIoError, AsyncIoOperation, AsyncIoResult, +}; use crate::error::{BlockError, BlockErrorKind, BlockResult}; use crate::sparse::{punch_hole, write_zeroes}; use crate::{SECTOR_SIZE, is_block_device}; pub struct RawFileAsyncAio { fd: RawFd, - ctx: aio::IoContext, - eventfd: EventFd, + data_io: AioDataIo, alignment: u64, - completion_list: VecDeque<(u64, i32)>, is_block_device: bool, } impl RawFileAsyncAio { pub fn new(fd: RawFd, queue_depth: u32) -> BlockResult { - let eventfd = - EventFd::new(libc::EFD_NONBLOCK).map_err(|e| BlockError::new(BlockErrorKind::Io, e))?; - let ctx = - aio::IoContext::new(queue_depth).map_err(|e| BlockError::new(BlockErrorKind::Io, e))?; + let data_io = + AioDataIo::new(queue_depth).map_err(|e| BlockError::new(BlockErrorKind::Io, e))?; let is_block_device = is_block_device(fd); Ok(RawFileAsyncAio { fd, - ctx, - eventfd, + data_io, alignment: SECTOR_SIZE, - completion_list: VecDeque::new(), is_block_device, }) } @@ -46,7 +42,7 @@ impl RawFileAsyncAio { impl AsyncIo for RawFileAsyncAio { fn notifier(&self) -> &EventFd { - &self.eventfd + self.data_io.notifier() } fn alignment(&self) -> u64 { @@ -59,23 +55,9 @@ impl AsyncIo for RawFileAsyncAio { iovecs: &[libc::iovec], user_data: u64, ) -> AsyncIoResult<()> { - let iocbs = [&mut aio::IoControlBlock { - aio_fildes: self.fd.as_raw_fd() as u32, - aio_lio_opcode: aio::IOCB_CMD_PREADV as u16, - aio_buf: iovecs.as_ptr() as u64, - aio_nbytes: iovecs.len() as u64, - aio_offset: offset, - aio_data: user_data, - aio_flags: aio::IOCB_FLAG_RESFD, - aio_resfd: self.eventfd.as_raw_fd() as u32, - ..Default::default() - }]; - let _ = self - .ctx - .submit(&iocbs[..]) - .map_err(AsyncIoError::ReadVectored)?; - - Ok(()) + self.data_io + .submit_borrowed_operation(self.fd, offset, true, iovecs, user_data) + .map_err(AsyncIoError::ReadVectored) } fn write_vectored( @@ -84,36 +66,27 @@ impl AsyncIo for RawFileAsyncAio { iovecs: &[libc::iovec], user_data: u64, ) -> AsyncIoResult<()> { - let iocbs = [&mut aio::IoControlBlock { - aio_fildes: self.fd.as_raw_fd() as u32, - aio_lio_opcode: aio::IOCB_CMD_PWRITEV as u16, - aio_buf: iovecs.as_ptr() as u64, - aio_nbytes: iovecs.len() as u64, - aio_offset: offset, - aio_data: user_data, - aio_flags: aio::IOCB_FLAG_RESFD, - aio_resfd: self.eventfd.as_raw_fd() as u32, - ..Default::default() - }]; - let _ = self - .ctx - .submit(&iocbs[..]) - .map_err(AsyncIoError::WriteVectored)?; + self.data_io + .submit_borrowed_operation(self.fd, offset, false, iovecs, user_data) + .map_err(AsyncIoError::WriteVectored) + } - Ok(()) + fn submit_data_operation(&mut self, op: AsyncIoOperation) -> AsyncIoResult<()> { + let is_read = op.is_read(); + self.data_io.submit_operation(self.fd, op).map_err(|e| { + if is_read { + AsyncIoError::ReadVectored(e) + } else { + AsyncIoError::WriteVectored(e) + } + }) } fn fsync(&mut self, user_data: Option) -> AsyncIoResult<()> { if let Some(user_data) = user_data { - let iocbs = [&mut aio::IoControlBlock { - aio_fildes: self.fd.as_raw_fd() as u32, - aio_lio_opcode: aio::IOCB_CMD_FSYNC as u16, - aio_data: user_data, - aio_flags: aio::IOCB_FLAG_RESFD, - aio_resfd: self.eventfd.as_raw_fd() as u32, - ..Default::default() - }]; - let _ = self.ctx.submit(&iocbs[..]).map_err(AsyncIoError::Fsync)?; + self.data_io + .submit_fsync(self.fd, user_data) + .map_err(AsyncIoError::Fsync)?; } else { // SAFETY: FFI call with a valid fd unsafe { libc::fsync(self.fd) }; @@ -122,17 +95,8 @@ impl AsyncIo for RawFileAsyncAio { Ok(()) } - fn next_completed_request(&mut self) -> Option<(u64, i32)> { - if self.completion_list.is_empty() { - // Drain pending AIO completions batched into the same queue. - let mut events = [aio::IoEvent::default(); 32]; - let rc = self.ctx.get_events(0, &mut events, None).unwrap(); - for event in &events[..rc] { - self.completion_list - .push_back((event.data, event.res as i32)); - } - } - 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<()> { @@ -141,8 +105,8 @@ impl AsyncIo for RawFileAsyncAio { // list, matching the pattern used by the sync backend (RawFileSync). punch_hole(self.fd, self.is_block_device, offset, length) .map_err(AsyncIoError::PunchHole)?; - 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(()) } @@ -151,8 +115,8 @@ impl AsyncIo for RawFileAsyncAio { // Same as punch_hole(). write_zeroes(self.fd, self.is_block_device, offset, length) .map_err(AsyncIoError::WriteZeroes)?; - 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(()) }