mirror of
https://github.com/cloud-hypervisor/cloud-hypervisor.git
synced 2026-08-05 02:19:16 +00:00
vhost_user: Adapt backends to let handle_event be immutable
Both blk, net and fs backends have been updated to avoid the requirement of having handle_event(&mut self). This will allow the backend crate to avoid taking a write lock onto the backend object, which will remove the potential contention point when multiple threads will be handling multiqueues. Signed-off-by: Sebastien Boeuf <sebastien.boeuf@intel.com>
This commit is contained in:
@@ -23,7 +23,7 @@ use std::io::{self};
|
||||
use std::net::Ipv4Addr;
|
||||
use std::os::unix::io::AsRawFd;
|
||||
use std::process;
|
||||
use std::sync::{Arc, RwLock};
|
||||
use std::sync::{Arc, Mutex, RwLock};
|
||||
use std::vec::Vec;
|
||||
use vhost_rs::vhost_user::message::*;
|
||||
use vhost_rs::vhost_user::Error as VhostUserError;
|
||||
@@ -92,7 +92,7 @@ impl std::convert::From<Error> for std::io::Error {
|
||||
}
|
||||
}
|
||||
|
||||
pub struct VhostUserNetBackend {
|
||||
struct VhostUserNetThread {
|
||||
mem: Option<GuestMemoryMmap>,
|
||||
vring_worker: Option<Arc<VringWorker>>,
|
||||
kill_evt: EventFd,
|
||||
@@ -101,12 +101,11 @@ pub struct VhostUserNetBackend {
|
||||
txs: Vec<TxVirtio>,
|
||||
rx_tap_listenings: Vec<bool>,
|
||||
num_queues: usize,
|
||||
queue_size: u16,
|
||||
}
|
||||
|
||||
impl VhostUserNetBackend {
|
||||
impl VhostUserNetThread {
|
||||
/// Create a new virtio network device with the given TAP interface.
|
||||
pub fn new_with_tap(taps: Vec<Tap>, num_queues: usize, queue_size: u16) -> Result<Self> {
|
||||
fn new_with_tap(taps: Vec<Tap>, num_queues: usize) -> Result<Self> {
|
||||
let mut taps_v: Vec<(Tap, usize)> = Vec::new();
|
||||
for (i, tap) in taps.iter().enumerate() {
|
||||
taps_v.push((tap.clone(), num_queues + i));
|
||||
@@ -124,7 +123,7 @@ impl VhostUserNetBackend {
|
||||
rx_tap_listenings.push(false);
|
||||
}
|
||||
|
||||
Ok(VhostUserNetBackend {
|
||||
Ok(VhostUserNetThread {
|
||||
mem: None,
|
||||
vring_worker: None,
|
||||
kill_evt: EventFd::new(EFD_NONBLOCK).map_err(Error::CreateKillEventFd)?,
|
||||
@@ -133,23 +132,21 @@ impl VhostUserNetBackend {
|
||||
txs,
|
||||
rx_tap_listenings,
|
||||
num_queues,
|
||||
queue_size,
|
||||
})
|
||||
}
|
||||
|
||||
/// Create a new virtio network device with the given IP address and
|
||||
/// netmask.
|
||||
pub fn new(
|
||||
fn new(
|
||||
ip_addr: Ipv4Addr,
|
||||
netmask: Ipv4Addr,
|
||||
num_queues: usize,
|
||||
queue_size: u16,
|
||||
ifname: Option<&str>,
|
||||
) -> Result<Self> {
|
||||
let taps = open_tap(ifname, Some(ip_addr), Some(netmask), num_queues / 2)
|
||||
.map_err(Error::OpenTap)?;
|
||||
|
||||
Self::new_with_tap(taps, num_queues, queue_size)
|
||||
Self::new_with_tap(taps, num_queues)
|
||||
}
|
||||
|
||||
// Copies a single frame from `self.rx.frame_buf` into the guest. Returns true
|
||||
@@ -252,6 +249,32 @@ impl VhostUserNetBackend {
|
||||
}
|
||||
}
|
||||
|
||||
pub struct VhostUserNetBackend {
|
||||
thread: Mutex<VhostUserNetThread>,
|
||||
num_queues: usize,
|
||||
queue_size: u16,
|
||||
}
|
||||
|
||||
impl VhostUserNetBackend {
|
||||
fn new(
|
||||
ip_addr: Ipv4Addr,
|
||||
netmask: Ipv4Addr,
|
||||
num_queues: usize,
|
||||
queue_size: u16,
|
||||
ifname: Option<&str>,
|
||||
) -> Result<Self> {
|
||||
let thread = Mutex::new(VhostUserNetThread::new(
|
||||
ip_addr, netmask, num_queues, ifname,
|
||||
)?);
|
||||
|
||||
Ok(VhostUserNetBackend {
|
||||
thread,
|
||||
num_queues,
|
||||
queue_size,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl VhostUserBackend for VhostUserNetBackend {
|
||||
fn num_queues(&self) -> usize {
|
||||
self.num_queues
|
||||
@@ -279,7 +302,7 @@ impl VhostUserBackend for VhostUserNetBackend {
|
||||
fn set_event_idx(&mut self, _enabled: bool) {}
|
||||
|
||||
fn update_memory(&mut self, mem: GuestMemoryMmap) -> VhostUserBackendResult<()> {
|
||||
self.mem = Some(mem);
|
||||
self.thread.lock().unwrap().mem = Some(mem);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -296,42 +319,43 @@ impl VhostUserBackend for VhostUserNetBackend {
|
||||
let tap_start_index = self.num_queues as u16;
|
||||
let tap_end_index = (self.num_queues + self.num_queues / 2 - 1) as u16;
|
||||
|
||||
let mut thread = self.thread.lock().unwrap();
|
||||
match device_event {
|
||||
x if ((x < self.num_queues as u16) && (x % 2 == 0)) => {
|
||||
let index = (x / 2) as usize;
|
||||
let mut vring = vrings[x as usize].write().unwrap();
|
||||
self.resume_rx(&mut vring, index)?;
|
||||
thread.resume_rx(&mut vring, index)?;
|
||||
|
||||
if !self.rx_tap_listenings[index] {
|
||||
self.vring_worker.as_ref().unwrap().register_listener(
|
||||
self.taps[index].0.as_raw_fd(),
|
||||
if !thread.rx_tap_listenings[index] {
|
||||
thread.vring_worker.as_ref().unwrap().register_listener(
|
||||
thread.taps[index].0.as_raw_fd(),
|
||||
epoll::Events::EPOLLIN,
|
||||
u64::try_from(self.taps[index].1).unwrap(),
|
||||
u64::try_from(thread.taps[index].1).unwrap(),
|
||||
)?;
|
||||
self.rx_tap_listenings[index] = true;
|
||||
thread.rx_tap_listenings[index] = true;
|
||||
}
|
||||
}
|
||||
x if ((x < self.num_queues as u16) && (x % 2 != 0)) => {
|
||||
x if ((x < thread.num_queues as u16) && (x % 2 != 0)) => {
|
||||
let index = ((x - 1) / 2) as usize;
|
||||
let mut vring = vrings[x as usize].write().unwrap();
|
||||
self.process_tx(&mut vring.mut_queue(), index)?;
|
||||
thread.process_tx(&mut vring.mut_queue(), index)?;
|
||||
}
|
||||
x if x >= tap_start_index && x <= tap_end_index => {
|
||||
let index = x as usize - self.num_queues;
|
||||
let mut vring = vrings[2 * index].write().unwrap();
|
||||
if self.rxs[index].deferred_frame
|
||||
if thread.rxs[index].deferred_frame
|
||||
// Process a deferred frame first if available. Don't read from tap again
|
||||
// until we manage to receive this deferred frame.
|
||||
{
|
||||
if self.rx_single_frame(&mut vring.mut_queue(), index)? {
|
||||
self.rxs[index].deferred_frame = false;
|
||||
self.process_rx(&mut vring, index)?;
|
||||
} else if self.rxs[index].deferred_irqs {
|
||||
self.rxs[index].deferred_irqs = false;
|
||||
if thread.rx_single_frame(&mut vring.mut_queue(), index)? {
|
||||
thread.rxs[index].deferred_frame = false;
|
||||
thread.process_rx(&mut vring, index)?;
|
||||
} else if thread.rxs[index].deferred_irqs {
|
||||
thread.rxs[index].deferred_irqs = false;
|
||||
vring.signal_used_queue()?;
|
||||
}
|
||||
} else {
|
||||
self.process_rx(&mut vring, index)?;
|
||||
thread.process_rx(&mut vring, index)?;
|
||||
}
|
||||
}
|
||||
_ => return Err(Error::HandleEventUnknownEvent.into()),
|
||||
@@ -343,7 +367,10 @@ impl VhostUserBackend for VhostUserNetBackend {
|
||||
fn exit_event(&self) -> Option<(EventFd, Option<u16>)> {
|
||||
let tap_end_index = (self.num_queues + self.num_queues / 2 - 1) as u16;
|
||||
let kill_index = tap_end_index + 1;
|
||||
Some((self.kill_evt.try_clone().unwrap(), Some(kill_index)))
|
||||
Some((
|
||||
self.thread.lock().unwrap().kill_evt.try_clone().unwrap(),
|
||||
Some(kill_index),
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -451,6 +478,9 @@ pub fn start_net_backend(backend_command: &str) {
|
||||
net_backend
|
||||
.write()
|
||||
.unwrap()
|
||||
.thread
|
||||
.lock()
|
||||
.unwrap()
|
||||
.set_vring_worker(Some(vring_worker));
|
||||
|
||||
if let Err(e) = net_daemon.start() {
|
||||
@@ -465,7 +495,15 @@ pub fn start_net_backend(backend_command: &str) {
|
||||
error!("Error from the main thread: {:?}", e);
|
||||
}
|
||||
|
||||
let kill_evt = &net_backend.write().unwrap().kill_evt;
|
||||
let kill_evt = net_backend
|
||||
.write()
|
||||
.unwrap()
|
||||
.thread
|
||||
.lock()
|
||||
.unwrap()
|
||||
.kill_evt
|
||||
.try_clone()
|
||||
.unwrap();
|
||||
if let Err(e) = kill_evt.write(1) {
|
||||
error!("Error shutting down worker thread: {:?}", e)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user