From 4dfc3dc7e034503f9a6d68c0a970904751551d47 Mon Sep 17 00:00:00 2001 From: Wenyu Huang Date: Fri, 29 Aug 2025 14:42:26 +0000 Subject: [PATCH 1/7] Change device_event in handle_event() from u16 to usize When register_event we pass in the value of u32/u64, so we don't need to convert it to u16 and then pass it to handle_event. Signed-off-by: Wenyu Huang --- vhost-user-backend/src/backend.rs | 12 ++++++------ vhost-user-backend/src/event_loop.rs | 12 ++++++------ vhost-user-backend/tests/vhost-user-server.rs | 2 +- 3 files changed, 13 insertions(+), 13 deletions(-) diff --git a/vhost-user-backend/src/backend.rs b/vhost-user-backend/src/backend.rs index cc7232cc..1190aecc 100644 --- a/vhost-user-backend/src/backend.rs +++ b/vhost-user-backend/src/backend.rs @@ -144,7 +144,7 @@ pub trait VhostUserBackend: Send + Sync { /// do with events happening on custom listeners. fn handle_event( &self, - device_event: u16, + device_event: usize, evset: EventSet, vrings: &[Self::Vring], thread_id: usize, @@ -295,7 +295,7 @@ pub trait VhostUserBackendMut: Send + Sync { /// do with events happening on custom listeners. fn handle_event( &mut self, - device_event: u16, + device_event: usize, evset: EventSet, vrings: &[Self::Vring], thread_id: usize, @@ -404,7 +404,7 @@ impl VhostUserBackend for Arc { fn handle_event( &self, - device_event: u16, + device_event: usize, evset: EventSet, vrings: &[Self::Vring], thread_id: usize, @@ -497,7 +497,7 @@ impl VhostUserBackend for Mutex { fn handle_event( &self, - device_event: u16, + device_event: usize, evset: EventSet, vrings: &[Self::Vring], thread_id: usize, @@ -593,7 +593,7 @@ impl VhostUserBackend for RwLock { fn handle_event( &self, - device_event: u16, + device_event: usize, evset: EventSet, vrings: &[Self::Vring], thread_id: usize, @@ -741,7 +741,7 @@ pub mod tests { fn handle_event( &mut self, - _device_event: u16, + _device_event: usize, _evset: EventSet, _vrings: &[VringRwLock], _thread_id: usize, diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index 09e91042..d17b2298 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -184,10 +184,10 @@ where } }; - let ev_type = event.data() as u16; + let ev_type = event.data(); // handle_event() returns true if an event is received from the exit event fd. - if self.handle_event(ev_type, evset)? { + if self.handle_event(ev_type as usize, evset)? { break 'epoll; } } @@ -196,13 +196,13 @@ where Ok(()) } - fn handle_event(&self, device_event: u16, evset: EventSet) -> VringEpollResult { - if self.exit_event_fd.is_some() && device_event as usize == self.backend.num_queues() { + fn handle_event(&self, device_event: usize, evset: EventSet) -> VringEpollResult { + if self.exit_event_fd.is_some() && device_event == self.backend.num_queues() { return Ok(true); } - if (device_event as usize) < self.vrings.len() { - let vring = &self.vrings[device_event as usize]; + if device_event < self.vrings.len() { + let vring = &self.vrings[device_event]; let enabled = vring .read_kick() .map_err(VringEpollError::HandleEventReadKick)?; diff --git a/vhost-user-backend/tests/vhost-user-server.rs b/vhost-user-backend/tests/vhost-user-server.rs index 3c8205a5..67c8ecfc 100644 --- a/vhost-user-backend/tests/vhost-user-server.rs +++ b/vhost-user-backend/tests/vhost-user-server.rs @@ -117,7 +117,7 @@ impl VhostUserBackendMut for MockVhostBackend { fn handle_event( &mut self, - _device_event: u16, + _device_event: usize, _evset: EventSet, _vrings: &[VringRwLock], _thread_id: usize, From 0204fe93fcce694f3ee1e5797367b5cb78fa3420 Mon Sep 17 00:00:00 2001 From: Wenyu Huang Date: Fri, 29 Aug 2025 14:21:16 +0000 Subject: [PATCH 2/7] Change u64 to usize in un/register_event We can safely use usize instead of u64, because normally we will not register a data that exceeds the size of usize Signed-off-by: Wenyu Huang --- vhost-user-backend/src/event_loop.rs | 26 ++++++++++++++++---------- vhost-user-backend/src/handler.rs | 4 ++-- 2 files changed, 18 insertions(+), 12 deletions(-) diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index d17b2298..bd0101af 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -116,9 +116,9 @@ where /// /// When this event is later triggered, the backend implementation of `handle_event` will be /// called. - pub fn register_listener(&self, fd: RawFd, ev_type: EventSet, data: u64) -> Result<()> { + pub fn register_listener(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { // `data` range [0...num_queues] is reserved for queues and exit event. - if data <= self.backend.num_queues() as u64 { + if data <= self.backend.num_queues() { Err(io::Error::from_raw_os_error(libc::EINVAL)) } else { self.register_event(fd, ev_type, data) @@ -129,23 +129,29 @@ where /// /// If the event is triggered after this function has been called, the event will be silently /// dropped. - pub fn unregister_listener(&self, fd: RawFd, ev_type: EventSet, data: u64) -> Result<()> { + pub fn unregister_listener(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { // `data` range [0...num_queues] is reserved for queues and exit event. - if data <= self.backend.num_queues() as u64 { + if data <= self.backend.num_queues() { Err(io::Error::from_raw_os_error(libc::EINVAL)) } else { self.unregister_event(fd, ev_type, data) } } - pub(crate) fn register_event(&self, fd: RawFd, ev_type: EventSet, data: u64) -> Result<()> { - self.epoll - .ctl(ControlOperation::Add, fd, EpollEvent::new(ev_type, data)) + pub(crate) fn register_event(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { + self.epoll.ctl( + ControlOperation::Add, + fd, + EpollEvent::new(ev_type, data as u64), + ) } - pub(crate) fn unregister_event(&self, fd: RawFd, ev_type: EventSet, data: u64) -> Result<()> { - self.epoll - .ctl(ControlOperation::Delete, fd, EpollEvent::new(ev_type, data)) + pub(crate) fn unregister_event(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { + self.epoll.ctl( + ControlOperation::Delete, + fd, + EpollEvent::new(ev_type, data as u64), + ) } /// Run the event poll loop to handle all pending events on registered fds. diff --git a/vhost-user-backend/src/handler.rs b/vhost-user-backend/src/handler.rs index e81d1f9d..475ee8f5 100644 --- a/vhost-user-backend/src/handler.rs +++ b/vhost-user-backend/src/handler.rs @@ -222,7 +222,7 @@ where if let Err(e) = self.handlers[thread_index].register_event( fd.as_raw_fd(), EventSet::IN, - u64::from(evt_idx), + evt_idx as usize, ) { if e.kind() != io::ErrorKind::AlreadyExists { // This could happen if we're asked by the frontend to enable an @@ -234,7 +234,7 @@ where let _ = self.handlers[thread_index].unregister_event( fd.as_raw_fd(), EventSet::IN, - u64::from(evt_idx), + evt_idx as usize, ); } break; From 76a8463e27bcd00fa18fce04e415c953d12fc43b Mon Sep 17 00:00:00 2001 From: Wenyu Huang Date: Mon, 7 Jul 2025 16:21:12 +0000 Subject: [PATCH 3/7] Use mio to replace Epoll Epoll is linux-specific. So we use mio, which is a cross-platform event notification, to replace Epoll. Signed-off-by: Wenyu Huang --- vhost-user-backend/CHANGELOG.md | 2 + vhost-user-backend/Cargo.toml | 1 + vhost-user-backend/README.md | 2 +- vhost-user-backend/src/backend.rs | 5 +- vhost-user-backend/src/event_loop.rs | 245 +++++++++++------- vhost-user-backend/src/handler.rs | 32 +-- vhost-user-backend/src/lib.rs | 12 +- vhost-user-backend/tests/vhost-user-server.rs | 3 +- 8 files changed, 173 insertions(+), 129 deletions(-) diff --git a/vhost-user-backend/CHANGELOG.md b/vhost-user-backend/CHANGELOG.md index 29295c60..f8b3081b 100644 --- a/vhost-user-backend/CHANGELOG.md +++ b/vhost-user-backend/CHANGELOG.md @@ -5,6 +5,8 @@ ### Added - [[#355]](https://github.com/rust-vmm/vhost/pull/355) Add an explicit shutdown handle for active daemon connections. ### Changed +- [[316](https://github.com/rust-vmm/vhost/pull/316)] Use mio to replace Epoll. Expose event_loop::EventSet. + ### Deprecated ### Fixed diff --git a/vhost-user-backend/Cargo.toml b/vhost-user-backend/Cargo.toml index 360b8b53..d11c78a3 100644 --- a/vhost-user-backend/Cargo.toml +++ b/vhost-user-backend/Cargo.toml @@ -20,6 +20,7 @@ postcopy = ["vhost/postcopy", "userfaultfd"] libc = "0.2.39" log = "0.4.17" userfaultfd = { version = "0.9.0", optional = true } +mio = { version = "1.0.4", features = ["os-poll", "os-ext"] } vhost = { path = "../vhost", version = "0.16.0", features = ["vhost-user-backend"] } virtio-bindings = { workspace = true } virtio-queue = { workspace = true } diff --git a/vhost-user-backend/README.md b/vhost-user-backend/README.md index 46b771ce..5c90ac14 100644 --- a/vhost-user-backend/README.md +++ b/vhost-user-backend/README.md @@ -22,7 +22,7 @@ where pub fn new(name: String, backend: S, atomic_mem: GuestMemoryAtomic>) -> Result; pub fn start(&mut self, listener: Listener) -> Result<()>; pub fn wait(&mut self) -> Result<()>; - pub fn get_epoll_handlers(&self) -> Vec>>; + pub fn get_poll_handlers(&self) -> Vec>>; } ``` diff --git a/vhost-user-backend/src/backend.rs b/vhost-user-backend/src/backend.rs index 1190aecc..0763dac6 100644 --- a/vhost-user-backend/src/backend.rs +++ b/vhost-user-backend/src/backend.rs @@ -29,7 +29,6 @@ use vhost::vhost_user::message::{ }; use vhost::vhost_user::Backend; use vm_memory::bitmap::Bitmap; -use vmm_sys_util::epoll::EventSet; use vmm_sys_util::event::{EventConsumer, EventNotifier}; use vhost::vhost_user::GpuBackend; @@ -37,6 +36,8 @@ use vhost::vhost_user::GpuBackend; use super::vring::VringT; use super::GM; +use crate::EventSet; + /// Trait with interior mutability for vhost user backend servers to implement concrete services. /// /// To support multi-threading and asynchronous IO, we enforce `Send + Sync` bound. @@ -828,7 +829,7 @@ pub mod tests { let vring = VringRwLock::new(mem, 0x1000).unwrap(); backend - .handle_event(0x1, EventSet::IN, &[vring], 0) + .handle_event(0x1, EventSet::Readable, &[vring], 0) .unwrap(); backend.reset_device(); diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index bd0101af..2b26c6c4 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -3,62 +3,103 @@ // // SPDX-License-Identifier: Apache-2.0 +use std::collections::HashSet; use std::fmt::{Display, Formatter}; use std::io::{self, Result}; use std::marker::PhantomData; use std::os::fd::IntoRawFd; use std::os::unix::io::{AsRawFd, RawFd}; +use std::sync::Mutex; -use vmm_sys_util::epoll::{ControlOperation, Epoll, EpollEvent, EventSet}; +use mio::event::Event; +use mio::unix::SourceFd; +use mio::{Events, Interest, Poll, Registry, Token}; use vmm_sys_util::event::EventNotifier; use super::backend::VhostUserBackend; use super::vring::VringT; -/// Errors related to vring epoll event handling. +/// Errors related to vring epoll/kqueue event handling. #[derive(Debug)] -pub enum VringEpollError { +pub enum VringPollError { /// Failed to create epoll file descriptor. - EpollCreateFd(io::Error), + PollerCreate(io::Error), /// Failed while waiting for events. - EpollWait(io::Error), + PollerWait(io::Error), /// Could not register exit event RegisterExitEvent(io::Error), /// Failed to read the event from kick EventFd. HandleEventReadKick(io::Error), /// Failed to handle the event from the backend. HandleEventBackendHandling(io::Error), + /// Failed to clone registry. + RegistryClone(io::Error), } -impl Display for VringEpollError { +impl Display for VringPollError { fn fmt(&self, f: &mut Formatter) -> std::fmt::Result { match self { - VringEpollError::EpollCreateFd(e) => write!(f, "cannot create epoll fd: {e}"), - VringEpollError::EpollWait(e) => write!(f, "failed to wait for epoll event: {e}"), - VringEpollError::RegisterExitEvent(e) => write!(f, "cannot register exit event: {e}"), - VringEpollError::HandleEventReadKick(e) => { + VringPollError::PollerCreate(e) => write!(f, "cannot create poller: {e}"), + VringPollError::PollerWait(e) => write!(f, "failed to wait for poller event: {e}"), + VringPollError::RegisterExitEvent(e) => write!(f, "cannot register exit event: {e}"), + VringPollError::HandleEventReadKick(e) => { write!(f, "cannot read vring kick event: {e}") } - VringEpollError::HandleEventBackendHandling(e) => { - write!(f, "failed to handle epoll event: {e}") + VringPollError::HandleEventBackendHandling(e) => { + write!(f, "failed to handle poll event: {e}") } + VringPollError::RegistryClone(e) => write!(f, "cannot clone poller's registry: {e}"), } } } -impl std::error::Error for VringEpollError {} +impl std::error::Error for VringPollError {} -/// Result of vring epoll operations. -pub type VringEpollResult = std::result::Result; +/// Result of vring epoll/kqueue operations. +pub type VringPollResult = std::result::Result; -/// Epoll event handler to manage and process epoll events for registered file descriptor. +#[derive(Debug, Clone, Copy)] +pub enum EventSet { + Readable, + Writable, + All, +} + +impl EventSet { + fn to_interest(self) -> Interest { + match self { + EventSet::Readable => Interest::READABLE, + EventSet::Writable => Interest::WRITABLE, + EventSet::All => Interest::READABLE | Interest::WRITABLE, + } + } +} + +fn event_to_event_set(evt: &Event) -> Option { + if evt.is_readable() && evt.is_writable() { + return Some(EventSet::All); + } + if evt.is_readable() { + return Some(EventSet::Readable); + } + if evt.is_writable() { + return Some(EventSet::Writable); + } + None +} + +/// Epoll/kqueue event handler to manage and process epoll/kqueue events for registered file descriptor. /// -/// The `VringEpollHandler` structure provides interfaces to: -/// - add file descriptors to be monitored by the epoll fd -/// - remove registered file descriptors from the epoll fd -/// - run the event loop to handle pending events on the epoll fd -pub struct VringEpollHandler { - epoll: Epoll, +/// The `VringPollHandler` structure provides interfaces to: +/// - add file descriptors to be monitored by the epoll/kqueue fd +/// - remove registered file descriptors from the epoll/kqueue fd +/// - run the event loop to handle pending events on the epoll/kqueue fd +pub struct VringPollHandler { + poller: Mutex, + registry: Registry, + // Record the registered fd. + // Because in mio, consecutive calls to register is unspecified behavior. + fd_set: Mutex>, backend: T, vrings: Vec, thread_id: usize, @@ -66,7 +107,7 @@ pub struct VringEpollHandler { phantom: PhantomData, } -impl VringEpollHandler { +impl VringPollHandler { /// Send `exit event` to break the event loop. pub fn send_exit_event(&self) { if let Some(eventfd) = self.exit_event_fd.as_ref() { @@ -75,35 +116,45 @@ impl VringEpollHandler { } } -impl VringEpollHandler +impl VringPollHandler where T: VhostUserBackend, { - /// Create a `VringEpollHandler` instance. + /// Create a `VringPollHandler` instance. pub(crate) fn new( backend: T, vrings: Vec, thread_id: usize, - ) -> VringEpollResult { - let epoll = Epoll::new().map_err(VringEpollError::EpollCreateFd)?; + ) -> VringPollResult { + let poller = Poll::new().map_err(VringPollError::PollerCreate)?; let exit_event_fd = backend.exit_event(thread_id); + let fd_set = Mutex::new(HashSet::new()); + let registry = poller + .registry() + .try_clone() + .map_err(VringPollError::RegistryClone)?; let exit_event_fd = if let Some((consumer, notifier)) = exit_event_fd { let id = backend.num_queues(); - epoll - .ctl( - ControlOperation::Add, - consumer.into_raw_fd(), - EpollEvent::new(EventSet::IN, id as u64), + + registry + .register( + &mut SourceFd(&consumer.as_raw_fd()), + Token(id), + Interest::READABLE, ) - .map_err(VringEpollError::RegisterExitEvent)?; + .map_err(VringPollError::RegisterExitEvent)?; + + fd_set.lock().unwrap().insert(consumer.into_raw_fd()); Some(notifier) } else { None }; - Ok(VringEpollHandler { - epoll, + Ok(VringPollHandler { + poller: Mutex::new(poller), + registry, + fd_set, backend, vrings, thread_id, @@ -112,7 +163,7 @@ where }) } - /// Register an event into the epoll fd. + /// Register an event into the epoll/kqueue fd. /// /// When this event is later triggered, the backend implementation of `handle_event` will be /// called. @@ -125,76 +176,67 @@ where } } - /// Unregister an event from the epoll fd. + /// Unregister an event from the epoll/kqueue fd. /// /// If the event is triggered after this function has been called, the event will be silently /// dropped. - pub fn unregister_listener(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { + pub fn unregister_listener(&self, fd: RawFd, data: usize) -> Result<()> { // `data` range [0...num_queues] is reserved for queues and exit event. if data <= self.backend.num_queues() { Err(io::Error::from_raw_os_error(libc::EINVAL)) } else { - self.unregister_event(fd, ev_type, data) + self.unregister_event(fd) } } pub(crate) fn register_event(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { - self.epoll.ctl( - ControlOperation::Add, - fd, - EpollEvent::new(ev_type, data as u64), - ) + let mut fd_set = self.fd_set.lock().unwrap(); + if fd_set.contains(&fd) { + return Err(io::Error::from_raw_os_error(libc::EEXIST)); + } + self.registry + .register(&mut SourceFd(&fd), Token(data), ev_type.to_interest()) + .map_err(std::io::Error::other)?; + fd_set.insert(fd); + Ok(()) } - pub(crate) fn unregister_event(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { - self.epoll.ctl( - ControlOperation::Delete, - fd, - EpollEvent::new(ev_type, data as u64), - ) + pub(crate) fn unregister_event(&self, fd: RawFd) -> Result<()> { + let mut fd_set = self.fd_set.lock().unwrap(); + if !fd_set.contains(&fd) { + return Err(io::Error::from_raw_os_error(libc::ENOENT)); + } + self.registry + .deregister(&mut SourceFd(&fd)) + .map_err(|e| std::io::Error::other(format!("Failed to deregister fd {fd}: {e}")))?; + fd_set.remove(&fd); + Ok(()) } /// Run the event poll loop to handle all pending events on registered fds. /// /// The event loop will be terminated once an event is received from the `exit event fd` /// associated with the backend. - pub(crate) fn run(&self) -> VringEpollResult<()> { - const EPOLL_EVENTS_LEN: usize = 100; - let mut events = vec![EpollEvent::new(EventSet::empty(), 0); EPOLL_EVENTS_LEN]; - - 'epoll: loop { - let num_events = match self.epoll.wait(-1, &mut events[..]) { - Ok(res) => res, - Err(e) => { - if e.kind() == io::ErrorKind::Interrupted { - // It's well defined from the epoll_wait() syscall - // documentation that the epoll loop can be interrupted - // before any of the requested events occurred or the - // timeout expired. In both those cases, epoll_wait() - // returns an error of type EINTR, but this should not - // be considered as a regular error. Instead it is more - // appropriate to retry, by calling into epoll_wait(). - continue; + pub(crate) fn run(&self) -> VringPollResult<()> { + const POLL_EVENTS_LEN: usize = 100; + + let mut events = Events::with_capacity(POLL_EVENTS_LEN); + 'poll: loop { + self.poller + .lock() + .unwrap() + .poll(&mut events, None) + .map_err(VringPollError::PollerWait)?; + + for event in &events { + let token = event.token(); + + if let Some(evt_set) = event_to_event_set(event) { + if self.handle_event(token.0, evt_set)? { + break 'poll; } - return Err(VringEpollError::EpollWait(e)); - } - }; - - for event in events.iter().take(num_events) { - let evset = match EventSet::from_bits(event.events) { - Some(evset) => evset, - None => { - let evbits = event.events; - println!("epoll: ignoring unknown event set: 0x{evbits:x}"); - continue; - } - }; - - let ev_type = event.data(); - - // handle_event() returns true if an event is received from the exit event fd. - if self.handle_event(ev_type as usize, evset)? { - break 'epoll; + } else { + println!("ignoring unknown event set: {:#x}", event.token().0); } } } @@ -202,7 +244,7 @@ where Ok(()) } - fn handle_event(&self, device_event: usize, evset: EventSet) -> VringEpollResult { + fn handle_event(&self, device_event: usize, evset: EventSet) -> VringPollResult { if self.exit_event_fd.is_some() && device_event == self.backend.num_queues() { return Ok(true); } @@ -211,7 +253,7 @@ where let vring = &self.vrings[device_event]; let enabled = vring .read_kick() - .map_err(VringEpollError::HandleEventReadKick)?; + .map_err(VringPollError::HandleEventReadKick)?; // If the vring is not enabled, it should not be processed. if !enabled { @@ -221,15 +263,15 @@ where self.backend .handle_event(device_event, evset, &self.vrings, self.thread_id) - .map_err(VringEpollError::HandleEventBackendHandling)?; + .map_err(VringPollError::HandleEventBackendHandling)?; Ok(false) } } -impl AsRawFd for VringEpollHandler { +impl AsRawFd for VringPollHandler { fn as_raw_fd(&self) -> RawFd { - self.epoll.as_raw_fd() + self.poller.lock().unwrap().as_raw_fd() } } @@ -243,40 +285,43 @@ mod tests { use vmm_sys_util::event::{new_event_consumer_and_notifier, EventFlag}; #[test] - fn test_vring_epoll_handler() { + fn test_vring_poll_handler() { let mem = GuestMemoryAtomic::new( GuestMemoryMmap::<()>::from_ranges(&[(GuestAddress(0x100000), 0x10000)]).unwrap(), ); let vring = VringRwLock::new(mem, 0x1000).unwrap(); let backend = Arc::new(Mutex::new(MockVhostBackend::new())); - let handler = VringEpollHandler::new(backend, vec![vring], 0x1).unwrap(); + let handler = VringPollHandler::new(backend, vec![vring], 0x1).unwrap(); let (consumer, _notifier) = new_event_consumer_and_notifier(EventFlag::empty()).unwrap(); handler - .register_listener(consumer.as_raw_fd(), EventSet::IN, 3) + .register_listener(consumer.as_raw_fd(), EventSet::Readable, 3) .unwrap(); // Register an already registered fd. handler - .register_listener(consumer.as_raw_fd(), EventSet::IN, 3) + .register_listener(consumer.as_raw_fd(), EventSet::Readable, 3) .unwrap_err(); // Register an invalid data. handler - .register_listener(consumer.as_raw_fd(), EventSet::IN, 1) + .register_listener(consumer.as_raw_fd(), EventSet::Readable, 1) .unwrap_err(); handler - .unregister_listener(consumer.as_raw_fd(), EventSet::IN, 3) + .unregister_listener(consumer.as_raw_fd(), 3) .unwrap(); // unregister an already unregistered fd. handler - .unregister_listener(consumer.as_raw_fd(), EventSet::IN, 3) + .unregister_listener(consumer.as_raw_fd(), 3) .unwrap_err(); // unregister an invalid data. handler - .unregister_listener(consumer.as_raw_fd(), EventSet::IN, 1) + .unregister_listener(consumer.as_raw_fd(), 1) .unwrap_err(); // Check we retrieve the correct file descriptor - assert_eq!(handler.as_raw_fd(), handler.epoll.as_raw_fd()); + assert_eq!( + handler.as_raw_fd(), + handler.poller.lock().unwrap().as_raw_fd() + ); } } diff --git a/vhost-user-backend/src/handler.rs b/vhost-user-backend/src/handler.rs index 475ee8f5..838e2e31 100644 --- a/vhost-user-backend/src/handler.rs +++ b/vhost-user-backend/src/handler.rs @@ -14,6 +14,7 @@ use std::sync::Arc; use std::thread; use crate::bitmap::{BitmapReplace, MemRegionBitmap, MmapLogReg}; +use crate::event_loop::EventSet; #[cfg(feature = "postcopy")] use userfaultfd::{Uffd, UffdBuilder}; use vhost::vhost_user::message::{ @@ -31,11 +32,10 @@ use virtio_bindings::bindings::virtio_ring::VIRTIO_RING_F_EVENT_IDX; use virtio_queue::{Error as VirtQueError, QueueT}; use vm_memory::mmap::NewBitmap; use vm_memory::{GuestAddress, GuestAddressSpace, GuestMemory, GuestMemoryMmap, GuestRegionMmap}; -use vmm_sys_util::epoll::EventSet; use super::backend::VhostUserBackend; -use super::event_loop::VringEpollHandler; -use super::event_loop::{VringEpollError, VringEpollResult}; +use super::event_loop::VringPollHandler; +use super::event_loop::{VringPollError, VringPollResult}; use super::vring::VringT; use super::GM; @@ -50,7 +50,7 @@ pub enum VhostUserHandlerError { /// Failed to create a `Vring`. CreateVring(VirtQueError), /// Failed to create vring worker. - CreateEpollHandler(VringEpollError), + CreatePollHandler(VringPollError), /// Failed to spawn vring worker. SpawnVringWorker(io::Error), /// Could not find the mapping from memory regions. @@ -63,8 +63,8 @@ impl std::fmt::Display for VhostUserHandlerError { VhostUserHandlerError::CreateVring(e) => { write!(f, "failed to create vring: {e}") } - VhostUserHandlerError::CreateEpollHandler(e) => { - write!(f, "failed to create vring epoll handler: {e}") + VhostUserHandlerError::CreatePollHandler(e) => { + write!(f, "failed to create vring poll handler: {e}") } VhostUserHandlerError::SpawnVringWorker(e) => { write!(f, "failed spawning the vring worker: {e}") @@ -90,7 +90,7 @@ struct AddrMapping { pub struct VhostUserHandler { backend: T, - handlers: Vec>>, + handlers: Vec>>, owned: bool, features_acked: bool, acked_features: u64, @@ -103,7 +103,7 @@ pub struct VhostUserHandler { vrings: Vec, #[cfg(feature = "postcopy")] uffd: Option, - worker_threads: Vec>>, + worker_threads: Vec>>, } // Ensure VhostUserHandler: Clone + Send + Sync + 'static. @@ -136,8 +136,8 @@ where } let handler = Arc::new( - VringEpollHandler::new(backend.clone(), thread_vrings, thread_id) - .map_err(VhostUserHandlerError::CreateEpollHandler)?, + VringPollHandler::new(backend.clone(), thread_vrings, thread_id) + .map_err(VhostUserHandlerError::CreatePollHandler)?, ); let handler2 = handler.clone(); let worker_thread = thread::Builder::new() @@ -191,7 +191,7 @@ impl VhostUserHandler where T: VhostUserBackend, { - pub(crate) fn get_epoll_handlers(&self) -> Vec>> { + pub(crate) fn get_poll_handlers(&self) -> Vec>> { self.handlers.clone() } @@ -208,7 +208,7 @@ where self.update_vring_registration(vring, index) } - /// Adds or removes the vring's kick fd to the epoll instance based on the vring status. + /// Adds or removes the vring's kick fd to the epoll/kqueue instance based on the vring status. /// Ensures that notifications are handled only while the vring is both started and enabled /// and that no notifications are lost. fn update_vring_registration(&self, vring: &T::Vring, index: u8) -> VhostUserResult<()> { @@ -221,7 +221,7 @@ where if vring_state.get_queue().ready() && vring_state.is_enabled() { if let Err(e) = self.handlers[thread_index].register_event( fd.as_raw_fd(), - EventSet::IN, + EventSet::Readable, evt_idx as usize, ) { if e.kind() != io::ErrorKind::AlreadyExists { @@ -231,11 +231,7 @@ where } } } else { - let _ = self.handlers[thread_index].unregister_event( - fd.as_raw_fd(), - EventSet::IN, - evt_idx as usize, - ); + let _ = self.handlers[thread_index].unregister_event(fd.as_raw_fd()); } break; } diff --git a/vhost-user-backend/src/lib.rs b/vhost-user-backend/src/lib.rs index da946322..ab384f6d 100644 --- a/vhost-user-backend/src/lib.rs +++ b/vhost-user-backend/src/lib.rs @@ -26,7 +26,7 @@ mod backend; pub use self::backend::{VhostUserBackend, VhostUserBackendMut}; mod event_loop; -pub use self::event_loop::VringEpollHandler; +pub use self::event_loop::{EventSet, VringPollHandler}; mod handler; pub use self::handler::VhostUserHandlerError; @@ -309,13 +309,13 @@ where } } - /// Retrieve the vring epoll handler. + /// Retrieve the vring poll handler. /// /// This is necessary to perform further actions like registering and unregistering some extra /// event file descriptors. - pub fn get_epoll_handlers(&self) -> Vec>> { + pub fn get_poll_handlers(&self) -> Vec>> { // Do not expect poisoned lock. - self.handler.lock().unwrap().get_epoll_handlers() + self.handler.lock().unwrap().get_poll_handlers() } } @@ -345,7 +345,7 @@ mod tests { let backend = Arc::new(Mutex::new(MockVhostBackend::new())); let mut daemon = VhostUserDaemon::new("test".to_owned(), backend, mem).unwrap(); - let handlers = daemon.get_epoll_handlers(); + let handlers = daemon.get_poll_handlers(); assert_eq!(handlers.len(), 2); let barrier = Arc::new(Barrier::new(2)); @@ -378,7 +378,7 @@ mod tests { let backend = Arc::new(Mutex::new(MockVhostBackend::new())); let mut daemon = VhostUserDaemon::new("test".to_owned(), backend, mem).unwrap(); - let handlers = daemon.get_epoll_handlers(); + let handlers = daemon.get_poll_handlers(); assert_eq!(handlers.len(), 2); let barrier = Arc::new(Barrier::new(2)); diff --git a/vhost-user-backend/tests/vhost-user-server.rs b/vhost-user-backend/tests/vhost-user-server.rs index 67c8ecfc..a78e5a84 100644 --- a/vhost-user-backend/tests/vhost-user-server.rs +++ b/vhost-user-backend/tests/vhost-user-server.rs @@ -13,11 +13,10 @@ use vhost::vhost_user::message::{ }; use vhost::vhost_user::{Backend, Frontend, Listener, VhostUserFrontend}; use vhost::{VhostBackend, VhostUserMemoryRegionInfo, VringConfigData}; -use vhost_user_backend::{VhostUserBackendMut, VhostUserDaemon, VringRwLock}; +use vhost_user_backend::{EventSet, VhostUserBackendMut, VhostUserDaemon, VringRwLock}; use vm_memory::{ FileOffset, GuestAddress, GuestAddressSpace, GuestMemory, GuestMemoryAtomic, GuestMemoryMmap, }; -use vmm_sys_util::epoll::EventSet; use vmm_sys_util::event::{ new_event_consumer_and_notifier, EventConsumer, EventFlag, EventNotifier, }; From 171b7f89ed84ac9199b27bae3e89232af9cc7297 Mon Sep 17 00:00:00 2001 From: EricMwangi Date: Wed, 1 Jul 2026 22:13:21 +0300 Subject: [PATCH 4/7] use log macro for unknown vring event change the logging od undefined events in vring are logged using `println!` to `log::warn!`. Signed-off-by: EricMwangi --- vhost-user-backend/src/event_loop.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index 2b26c6c4..d5d9e734 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -236,7 +236,7 @@ where break 'poll; } } else { - println!("ignoring unknown event set: {:#x}", event.token().0); + log::warn!("ignoring unknown event set: {:#x}", event.token().0); } } } From 794b34068cc5b23bcb076e49466f7aa021128a6e Mon Sep 17 00:00:00 2001 From: EricMwangi Date: Wed, 1 Jul 2026 23:51:40 +0300 Subject: [PATCH 5/7] rename un/register_listener parameter to queue_idx Rename the generic third paremeter in un/register_listener to `queue_idx`. This explicitly describes the function of the parameter. Signed-off-by: EricMwangi --- vhost-user-backend/src/event_loop.rs | 14 +++++++------- 1 file changed, 7 insertions(+), 7 deletions(-) diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index d5d9e734..a4f8bfc1 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -167,12 +167,12 @@ where /// /// When this event is later triggered, the backend implementation of `handle_event` will be /// called. - pub fn register_listener(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { + pub fn register_listener(&self, fd: RawFd, ev_type: EventSet, queue_idx: usize) -> Result<()> { // `data` range [0...num_queues] is reserved for queues and exit event. - if data <= self.backend.num_queues() { + if queue_idx <= self.backend.num_queues() { Err(io::Error::from_raw_os_error(libc::EINVAL)) } else { - self.register_event(fd, ev_type, data) + self.register_event(fd, ev_type, queue_idx) } } @@ -180,22 +180,22 @@ where /// /// If the event is triggered after this function has been called, the event will be silently /// dropped. - pub fn unregister_listener(&self, fd: RawFd, data: usize) -> Result<()> { + pub fn unregister_listener(&self, fd: RawFd, queue_idx: usize) -> Result<()> { // `data` range [0...num_queues] is reserved for queues and exit event. - if data <= self.backend.num_queues() { + if queue_idx <= self.backend.num_queues() { Err(io::Error::from_raw_os_error(libc::EINVAL)) } else { self.unregister_event(fd) } } - pub(crate) fn register_event(&self, fd: RawFd, ev_type: EventSet, data: usize) -> Result<()> { + pub(crate) fn register_event(&self, fd: RawFd, ev_type: EventSet, queue_idx: usize) -> Result<()> { let mut fd_set = self.fd_set.lock().unwrap(); if fd_set.contains(&fd) { return Err(io::Error::from_raw_os_error(libc::EEXIST)); } self.registry - .register(&mut SourceFd(&fd), Token(data), ev_type.to_interest()) + .register(&mut SourceFd(&fd), Token(queue_idx), ev_type.to_interest()) .map_err(std::io::Error::other)?; fd_set.insert(fd); Ok(()) From cacb9fc8e996fdaeb950ebc79eb039a59703b3c0 Mon Sep 17 00:00:00 2001 From: EricMwangi Date: Thu, 2 Jul 2026 02:48:01 +0300 Subject: [PATCH 6/7] Modify event_to_event_set to exhaust event states Add is_error,is_read_closed and is_write_closed to exhaust states of event. This provides better clarity into the state of events. Signed-off-by: EricMwangi --- vhost-user-backend/src/event_loop.rs | 42 ++++++++++++++++++---------- 1 file changed, 27 insertions(+), 15 deletions(-) diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index a4f8bfc1..8d75923b 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -75,17 +75,22 @@ impl EventSet { } } -fn event_to_event_set(evt: &Event) -> Option { +fn event_to_event_set(evt: &Event) -> Result { if evt.is_readable() && evt.is_writable() { - return Some(EventSet::All); + Ok(EventSet::All) + } else if evt.is_readable() { + Ok(EventSet::Readable) + } else if evt.is_writable() { + Ok(EventSet::Writable) + } else if evt.is_read_closed() { + Err(io::Error::other("Event Err is read closed")) + } else if evt.is_write_closed() { + Err(io::Error::other("Event Err is write closed")) + } else if evt.is_error() { + Err(io::Error::other("Event Epoll Error")) + } else { + Err(io::Error::other("Unknown or unhandled event state")) } - if evt.is_readable() { - return Some(EventSet::Readable); - } - if evt.is_writable() { - return Some(EventSet::Writable); - } - None } /// Epoll/kqueue event handler to manage and process epoll/kqueue events for registered file descriptor. @@ -231,13 +236,20 @@ where for event in &events { let token = event.token(); - if let Some(evt_set) = event_to_event_set(event) { - if self.handle_event(token.0, evt_set)? { - break 'poll; + match event_to_event_set(event) { + Ok(evt_set) => { + if self.handle_event(token.0, evt_set)? { + break 'poll; + } + } + Err(err) => { + log::warn!( + "Ignoring unknown event set for token {:#x} error: {:?}", + token.0, + err + ); } - } else { - log::warn!("ignoring unknown event set: {:#x}", event.token().0); - } + }; } } From e0b315a1463ed2ce79688fe1a1fc27f2eedb93ab Mon Sep 17 00:00:00 2001 From: EricMwangi Date: Fri, 3 Jul 2026 08:44:52 +0300 Subject: [PATCH 7/7] refactor: remove mutex wrapper from poller Modify `VringPollHandler` to pass `&mut Poll` to `run`, removing the need for a Mutex wrapper around Poller. This simplifies compliance with the `Send` + `Sync` requirements of `VringPollHandler` and removes false implications of multithreaded safety. Signed-off-by: EricMwangi --- vhost-user-backend/src/event_loop.rs | 23 +++++++++-------------- vhost-user-backend/src/handler.rs | 8 ++++++-- 2 files changed, 15 insertions(+), 16 deletions(-) diff --git a/vhost-user-backend/src/event_loop.rs b/vhost-user-backend/src/event_loop.rs index 8d75923b..1d497486 100644 --- a/vhost-user-backend/src/event_loop.rs +++ b/vhost-user-backend/src/event_loop.rs @@ -100,7 +100,6 @@ fn event_to_event_set(evt: &Event) -> Result { /// - remove registered file descriptors from the epoll/kqueue fd /// - run the event loop to handle pending events on the epoll/kqueue fd pub struct VringPollHandler { - poller: Mutex, registry: Registry, // Record the registered fd. // Because in mio, consecutive calls to register is unspecified behavior. @@ -130,8 +129,8 @@ where backend: T, vrings: Vec, thread_id: usize, + poller: &Poll, ) -> VringPollResult { - let poller = Poll::new().map_err(VringPollError::PollerCreate)?; let exit_event_fd = backend.exit_event(thread_id); let fd_set = Mutex::new(HashSet::new()); @@ -157,7 +156,6 @@ where }; Ok(VringPollHandler { - poller: Mutex::new(poller), registry, fd_set, backend, @@ -222,14 +220,12 @@ where /// /// The event loop will be terminated once an event is received from the `exit event fd` /// associated with the backend. - pub(crate) fn run(&self) -> VringPollResult<()> { + pub(crate) fn run(&self, mut poller: Poll) -> VringPollResult<()> { const POLL_EVENTS_LEN: usize = 100; let mut events = Events::with_capacity(POLL_EVENTS_LEN); 'poll: loop { - self.poller - .lock() - .unwrap() + poller .poll(&mut events, None) .map_err(VringPollError::PollerWait)?; @@ -283,7 +279,7 @@ where impl AsRawFd for VringPollHandler { fn as_raw_fd(&self) -> RawFd { - self.poller.lock().unwrap().as_raw_fd() + self.registry.as_raw_fd() } } @@ -304,7 +300,8 @@ mod tests { let vring = VringRwLock::new(mem, 0x1000).unwrap(); let backend = Arc::new(Mutex::new(MockVhostBackend::new())); - let handler = VringPollHandler::new(backend, vec![vring], 0x1).unwrap(); + let poller = Poll::new().unwrap(); + let handler = VringPollHandler::new(backend, vec![vring], 0x1, &poller).unwrap(); let (consumer, _notifier) = new_event_consumer_and_notifier(EventFlag::empty()).unwrap(); handler @@ -330,10 +327,8 @@ mod tests { handler .unregister_listener(consumer.as_raw_fd(), 1) .unwrap_err(); - // Check we retrieve the correct file descriptor - assert_eq!( - handler.as_raw_fd(), - handler.poller.lock().unwrap().as_raw_fd() - ); + // Validate registry fd + assert_ne!(handler.as_raw_fd(), -1); + assert_ne!(handler.as_raw_fd(), poller.as_raw_fd()); } } diff --git a/vhost-user-backend/src/handler.rs b/vhost-user-backend/src/handler.rs index 838e2e31..6dbba547 100644 --- a/vhost-user-backend/src/handler.rs +++ b/vhost-user-backend/src/handler.rs @@ -15,6 +15,7 @@ use std::thread; use crate::bitmap::{BitmapReplace, MemRegionBitmap, MmapLogReg}; use crate::event_loop::EventSet; +use mio::Poll; #[cfg(feature = "postcopy")] use userfaultfd::{Uffd, UffdBuilder}; use vhost::vhost_user::message::{ @@ -135,14 +136,17 @@ where } } + let poller = Poll::new() + .map_err(VringPollError::PollerCreate) + .map_err(VhostUserHandlerError::CreatePollHandler)?; let handler = Arc::new( - VringPollHandler::new(backend.clone(), thread_vrings, thread_id) + VringPollHandler::new(backend.clone(), thread_vrings, thread_id, &poller) .map_err(VhostUserHandlerError::CreatePollHandler)?, ); let handler2 = handler.clone(); let worker_thread = thread::Builder::new() .name("vring_worker".to_string()) - .spawn(move || handler2.run()) + .spawn(move || handler2.run(poller)) .map_err(VhostUserHandlerError::SpawnVringWorker)?; handlers.push(handler);