virtio-devices: vsock: Adapt used handling to match other devices

Follow the same pattern as other virtio devices using a bool to check if
it needs notification and propagating its own Error enum.

Sadly this does still use `anyhow!()` but this does match with the
behaviour of the other devices in their implementations.

As a side effect we can now remove two errors from the top-level Error
enum in virtio-devices as these were only used by this module and those
errors had mangled descriptions.

Signed-off-by: Rob Bradford <rbradford@meta.com>
This commit is contained in:
Rob Bradford
2026-05-04 15:49:55 +01:00
parent 7d8986aad0
commit 63d6beb170
2 changed files with 54 additions and 24 deletions

View File

@@ -117,10 +117,6 @@ pub enum Error {
SetShmRegionsNotSupported, SetShmRegionsNotSupported,
#[error("Failed to process net queue")] #[error("Failed to process net queue")]
NetQueuePair(#[source] ::net_util::NetQueuePairError), NetQueuePair(#[source] ::net_util::NetQueuePairError),
#[error("Failed to ")]
QueueAddUsed(#[source] virtio_queue::Error),
#[error("Failed to ")]
QueueIterator(#[source] virtio_queue::Error),
} }
#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq)] #[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq)]

View File

@@ -20,6 +20,7 @@ use event_monitor::event;
use log::{debug, error, info, warn}; use log::{debug, error, info, warn};
use seccompiler::SeccompAction; use seccompiler::SeccompAction;
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use thiserror::Error;
use virtio_queue::{Queue, QueueOwnedT, QueueT}; use virtio_queue::{Queue, QueueOwnedT, QueueT};
use vm_memory::{GuestAddressSpace, GuestMemoryAtomic}; use vm_memory::{GuestAddressSpace, GuestMemoryAtomic};
use vm_migration::{Migratable, MigratableError, Pausable, Snapshot, Snapshottable, Transportable}; use vm_migration::{Migratable, MigratableError, Pausable, Snapshot, Snapshottable, Transportable};
@@ -68,6 +69,12 @@ pub const EVT_QUEUE_EVENT: u16 = EPOLL_HELPER_EVENT_LAST + 3;
// Notification coming from the backend. // Notification coming from the backend.
pub const BACKEND_EVENT: u16 = EPOLL_HELPER_EVENT_LAST + 4; pub const BACKEND_EVENT: u16 = EPOLL_HELPER_EVENT_LAST + 4;
#[derive(Error, Debug)]
enum Error {
#[error("Failed adding used index")]
QueueAddUsed(#[source] virtio_queue::Error),
}
/// The `VsockEpollHandler` implements the runtime logic of our vsock device: /// The `VsockEpollHandler` implements the runtime logic of our vsock device:
/// 1. Respond to TX queue events by wrapping virtio buffers into `VsockPacket`s, then sending those /// 1. Respond to TX queue events by wrapping virtio buffers into `VsockPacket`s, then sending those
/// packets to the `VsockBackend`; /// packets to the `VsockBackend`;
@@ -118,7 +125,7 @@ where
/// Walk the driver-provided RX queue buffers and attempt to fill them up with any data that we /// Walk the driver-provided RX queue buffers and attempt to fill them up with any data that we
/// have pending. /// have pending.
/// ///
fn process_rx(&mut self) -> result::Result<(), DeviceError> { fn process_rx(&mut self) -> result::Result<bool, Error> {
debug!("vsock: epoll_handler::process_rx()"); debug!("vsock: epoll_handler::process_rx()");
let mut used_descs = false; let mut used_descs = false;
@@ -155,21 +162,17 @@ where
self.queues[0] self.queues[0]
.add_used(desc_chain.memory(), desc_chain.head_index(), used_len) .add_used(desc_chain.memory(), desc_chain.head_index(), used_len)
.map_err(DeviceError::QueueAddUsed)?; .map_err(Error::QueueAddUsed)?;
used_descs = true; used_descs = true;
} }
if used_descs { Ok(used_descs)
self.signal_used_queue(0)
} else {
Ok(())
}
} }
/// Walk the driver-provided TX queue buffers, package them up as vsock packets, and send them to /// Walk the driver-provided TX queue buffers, package them up as vsock packets, and send them to
/// the backend for processing. /// the backend for processing.
/// ///
fn process_tx(&mut self) -> result::Result<(), DeviceError> { fn process_tx(&mut self) -> result::Result<bool, Error> {
debug!("vsock: epoll_handler::process_tx()"); debug!("vsock: epoll_handler::process_tx()");
let mut used_descs = false; let mut used_descs = false;
@@ -184,7 +187,7 @@ where
error!("vsock: error reading TX packet: {e:?}"); error!("vsock: error reading TX packet: {e:?}");
self.queues[1] self.queues[1]
.add_used(desc_chain.memory(), desc_chain.head_index(), 0) .add_used(desc_chain.memory(), desc_chain.head_index(), 0)
.map_err(DeviceError::QueueAddUsed)?; .map_err(Error::QueueAddUsed)?;
used_descs = true; used_descs = true;
continue; continue;
} }
@@ -197,15 +200,11 @@ where
self.queues[1] self.queues[1]
.add_used(desc_chain.memory(), desc_chain.head_index(), 0) .add_used(desc_chain.memory(), desc_chain.head_index(), 0)
.map_err(DeviceError::QueueAddUsed)?; .map_err(Error::QueueAddUsed)?;
used_descs = true; used_descs = true;
} }
if used_descs { Ok(used_descs)
self.signal_used_queue(1)
} else {
Ok(())
}
} }
fn run( fn run(
@@ -250,9 +249,16 @@ where
EpollHelperError::HandleEvent(anyhow!("Failed to get RX queue event: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to get RX queue event: {e:?}"))
})?; })?;
if self.backend.read().unwrap().has_pending_rx() { if self.backend.read().unwrap().has_pending_rx() {
self.process_rx().map_err(|e| { let needs_notification = self.process_rx().map_err(|e| {
EpollHelperError::HandleEvent(anyhow!("Failed to process RX queue: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to process RX queue: {e:?}"))
})?; })?;
if needs_notification {
self.signal_used_queue(0).map_err(|e| {
EpollHelperError::HandleEvent(anyhow!(
"Failed to signal used RX queue: {e:?}"
))
})?;
}
} }
} }
TX_QUEUE_EVENT => { TX_QUEUE_EVENT => {
@@ -261,17 +267,31 @@ where
EpollHelperError::HandleEvent(anyhow!("Failed to get TX queue event: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to get TX queue event: {e:?}"))
})?; })?;
self.process_tx().map_err(|e| { let needs_notification = self.process_tx().map_err(|e| {
EpollHelperError::HandleEvent(anyhow!("Failed to process TX queue: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to process TX queue: {e:?}"))
})?; })?;
if needs_notification {
self.signal_used_queue(1).map_err(|e| {
EpollHelperError::HandleEvent(anyhow!(
"Failed to signal used TX queue: {e:?}"
))
})?;
}
// The backend may have queued up responses to the packets we sent during TX queue // The backend may have queued up responses to the packets we sent during TX queue
// processing. If that happened, we need to fetch those responses and place them // processing. If that happened, we need to fetch those responses and place them
// into RX buffers. // into RX buffers.
if self.backend.read().unwrap().has_pending_rx() { if self.backend.read().unwrap().has_pending_rx() {
self.process_rx().map_err(|e| { let needs_notification = self.process_rx().map_err(|e| {
EpollHelperError::HandleEvent(anyhow!("Failed to process RX queue: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to process RX queue: {e:?}"))
})?; })?;
if needs_notification {
self.signal_used_queue(0).map_err(|e| {
EpollHelperError::HandleEvent(anyhow!(
"Failed to signal used RX queue: {e:?}"
))
})?;
}
} }
} }
EVT_QUEUE_EVENT => { EVT_QUEUE_EVENT => {
@@ -288,13 +308,27 @@ where
// In particular, if `self.backend.send_pkt()` halted the TX queue processing (by // In particular, if `self.backend.send_pkt()` halted the TX queue processing (by
// returning an error) at some point in the past, now is the time to try walking the // returning an error) at some point in the past, now is the time to try walking the
// TX queue again. // TX queue again.
self.process_tx().map_err(|e| { let needs_notification = self.process_tx().map_err(|e| {
EpollHelperError::HandleEvent(anyhow!("Failed to process TX queue: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to process TX queue: {e:?}"))
})?; })?;
if needs_notification {
self.signal_used_queue(1).map_err(|e| {
EpollHelperError::HandleEvent(anyhow!(
"Failed to signal used TX queue: {e:?}"
))
})?;
}
if self.backend.read().unwrap().has_pending_rx() { if self.backend.read().unwrap().has_pending_rx() {
self.process_rx().map_err(|e| { let needs_notification = self.process_rx().map_err(|e| {
EpollHelperError::HandleEvent(anyhow!("Failed to process RX queue: {e:?}")) EpollHelperError::HandleEvent(anyhow!("Failed to process RX queue: {e:?}"))
})?; })?;
if needs_notification {
self.signal_used_queue(0).map_err(|e| {
EpollHelperError::HandleEvent(anyhow!(
"Failed to signal used RX queue: {e:?}"
))
})?;
}
} }
} }
_ => { _ => {