From f73f8d66037d2d80826b435a6d18b37b6d0c4285 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Tue, 29 Sep 2026 08:41:41 -0300 Subject: [PATCH 1/2] vmm: wait for fast Threaded block completions in the queue event A request served by the Threaded engine takes two wake-ups of the event loop, one to submit it and one for its completion event, where the Sync engine takes one. On a fast backing file the wake-ups cost more than the IO: a guest doing buffered writes, each followed by fdatasync, commits about 9% slower than with Sync. After handing requests to the worker, wait for their completions, for up to 100us, and complete them in the same queue event. The worker raises no completion event for what the device gets that way: the device tells it how many completions it has received, and when it is waiting for more. The wait is only worth it while the backing file answers within that time. The first request that does not turns it off, and completions go through the completion event again, until one arrives fast enough. The event loop is never held for more than the 100us. With a guest doing buffered 32 KiB writes with fdatasync on a tmpfs-backed drive, cache_type Writeback, and nothing else going on: commits/s event loop wake-ups per commit Sync 15060-15301 2.0 Threaded, before 13852-13967 4.0 Threaded, after 14872-15163 2.0 On a backing file that adds 2 ms per request nothing changes: the wait is off, and other devices are served as before. --- .../src/devices/virtio/block/virtio/device.rs | 6 +- .../virtio/block/virtio/io/threaded_io.rs | 195 +++++++++++++++++- .../devices/virtio/block/virtio/threaded.rs | 49 +++++ 3 files changed, 240 insertions(+), 10 deletions(-) diff --git a/src/vmm/src/devices/virtio/block/virtio/device.rs b/src/vmm/src/devices/virtio/block/virtio/device.rs index e91d9af8820..91f104726c1 100644 --- a/src/vmm/src/devices/virtio/block/virtio/device.rs +++ b/src/vmm/src/devices/virtio/block/virtio/device.rs @@ -469,11 +469,7 @@ impl VirtioBlock { { error!("BlockError submitting pending block requests: {:?}", err); } - if let FileEngine::Threaded(ref mut engine) = self.disk.file_engine - && let Err(err) = engine.kick() - { - error!("BlockError submitting pending block requests: {:?}", err); - } + self.threaded_kick(); if !used_any { self.metrics.no_avail_buffer.inc(); diff --git a/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs b/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs index c5c94cb96dd..bfdc6a6f193 100644 --- a/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs +++ b/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs @@ -12,10 +12,16 @@ //! Requests are vectored: one request carries up to [`THREADED_SEG_MAX`] guest buffers and is //! served by a single `preadv`/`pwritev`. They are handed to the worker in batches, one per //! [`ThreadedFileEngine::kick`], so that a pass over the virtqueue wakes the worker once. +//! +//! While the backing file is fast, the submitter waits for the completions right after the kick, +//! for up to [`POLL_BUDGET`], see [`ThreadedFileEngine::kick_and_poll`]. A request is then submitted and +//! completed in one wake-up of the event loop, as with the sync engine, where the eventfd would +//! take a second one. use std::collections::VecDeque; use std::fs::File; use std::os::unix::io::AsRawFd; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, OnceLock, mpsc}; use std::thread; use std::time::{Duration, Instant}; @@ -41,6 +47,10 @@ pub const THREADED_IO_MAX_IN_FLIGHT: usize = 128; /// request takes this many entries of the 256-entry queue, plus two. pub const THREADED_SEG_MAX: u32 = 32; +/// How long the submitter waits for completions after a kick, at most. This is time the event +/// loop serves nothing else, so it has to stay far below what any other device would notice. +pub const POLL_BUDGET: Duration = Duration::from_micros(100); + /// A guest buffer: its address and length. pub type Segment = (GuestAddress, u32); @@ -86,6 +96,15 @@ pub enum ThreadedIoError { WorkerGone, } +/// What the submitter tells the worker about the completions it needs no signal for. +#[derive(Debug, Default)] +struct Received { + // Set while the submitter is polling for completions. + polling: AtomicBool, + // How many completions the submitter has received so far. + count: AtomicU64, +} + /// A finished request, as reported by the worker thread. #[derive(Debug)] pub struct ThreadedCompletion { @@ -250,14 +269,19 @@ struct CompletionSignal { evt: EventFd, pending: bool, last: Instant, + // How many completions were sent so far. + sent: u64, + received: Arc, } impl CompletionSignal { - fn new(evt: EventFd) -> Self { + fn new(evt: EventFd, received: Arc) -> Self { CompletionSignal { evt, pending: false, last: Instant::now(), + sent: 0, + received, } } @@ -273,10 +297,18 @@ impl CompletionSignal { if !self.pending { return; } + self.pending = false; + // No signal for completions the submitter already has, or is about to get: they were + // sent before this check, and the submitter looks for completions once more after it + // stops polling, so it cannot miss them. + if self.received.count.load(Ordering::SeqCst) >= self.sent + || self.received.polling.load(Ordering::SeqCst) + { + return; + } if let Err(err) = self.evt.write(1) { error!("Failed to signal block IO completion: {:?}", err); } - self.pending = false; self.last = Instant::now(); } } @@ -286,6 +318,7 @@ fn run_worker( batches: mpsc::Receiver>, completions: mpsc::Sender, completion_evt: EventFd, + received: Arc, ) { if let Some(filter) = worker_seccomp_filter() && let Err(err) = crate::seccomp::apply_filter(filter) @@ -293,7 +326,7 @@ fn run_worker( panic!("Failed to set the requested seccomp filters on the block IO worker: {err}"); } - let mut signal = CompletionSignal::new(completion_evt); + let mut signal = CompletionSignal::new(completion_evt, received); let mut queue = VecDeque::new(); loop { let op = match queue.pop_front() { @@ -331,6 +364,7 @@ fn run_worker( if completions.send(completion).is_err() { break; } + signal.sent += 1; signal.pending = true; // Hold the signal back while more requests are queued, so that the device handles a @@ -358,8 +392,15 @@ pub struct ThreadedFileEngine { // Requests pushed since the last kick. batch: Vec, completions: mpsc::Receiver, + // Completions received while polling, not popped yet. + ready: VecDeque, completion_evt: EventFd, in_flight: usize, + received: Arc, + // Whether the backing file has been answering within the poll budget. + poll_pays: bool, + poll_budget: Option, + last_kick: Instant, worker: Option>, } @@ -372,10 +413,20 @@ impl ThreadedFileEngine { let worker_file = file.try_clone().map_err(ThreadedIoError::FileClone)?; let (batches, worker_batches) = mpsc::channel(); let (worker_completions, completions) = mpsc::channel(); + let received = Arc::new(Received::default()); + let worker_received = received.clone(); let worker = thread::Builder::new() .name("fc_blk_io".to_string()) - .spawn(move || run_worker(worker_file, worker_batches, worker_completions, worker_evt)) + .spawn(move || { + run_worker( + worker_file, + worker_batches, + worker_completions, + worker_evt, + worker_received, + ) + }) .map_err(ThreadedIoError::Spawn)?; Ok(ThreadedFileEngine { @@ -383,8 +434,14 @@ impl ThreadedFileEngine { batches, batch: Vec::new(), completions, + ready: VecDeque::new(), completion_evt, in_flight: 0, + received, + // Tests expect completions to take the completion event, unless they ask for this. + poll_pays: !cfg!(test), + poll_budget: (!cfg!(test)).then_some(POLL_BUDGET), + last_kick: Instant::now(), worker: Some(worker), }) } @@ -411,6 +468,7 @@ impl ThreadedFileEngine { if self.batch.is_empty() { return Ok(()); } + self.last_kick = Instant::now(); self.batches .send(std::mem::take(&mut self.batch)) .map_err(|_| ThreadedIoError::WorkerGone) @@ -493,11 +551,77 @@ impl ThreadedFileEngine { /// Pop a finished request, if there is one. pub fn pop(&mut self) -> Option { - let completion = self.completions.try_recv().ok()?; + let completion = match self.ready.pop_front() { + Some(completion) => completion, + None => { + let completion = self.receive()?; + // Completions that take the slow path tell how fast the backing file is. + if self.in_flight == 1 + && let Some(budget) = self.poll_budget + { + self.poll_pays = self.last_kick.elapsed() < budget; + } + completion + } + }; self.in_flight -= 1; Some(completion) } + /// Take a completion from the worker, and let it know. + fn receive(&mut self) -> Option { + let completion = self.completions.try_recv().ok()?; + self.received.count.fetch_add(1, Ordering::SeqCst); + Some(completion) + } + + /// Turn waiting for completions after a kick on or off, see [`Self::kick_and_poll`]. + pub fn set_poll_budget(&mut self, budget: Option) { + self.poll_budget = budget; + self.poll_pays = budget.is_some(); + } + + /// Hand the requests pushed since the last kick to the worker, then wait for the completions + /// of everything in flight, for up to the poll budget, and only while the backing file has + /// been answering within it. Returns whether there are completions to pop. + /// + /// The worker does not signal the completions it finishes meanwhile. + pub fn kick_and_poll(&mut self) -> Result { + let budget = match self.poll_budget { + Some(budget) if self.poll_pays && self.in_flight > self.ready.len() => budget, + _ => { + self.kick()?; + return Ok(!self.ready.is_empty()); + } + }; + + // Before the worker gets the requests, so that it signals none of them. + self.received.polling.store(true, Ordering::SeqCst); + let kicked = self.kick(); + let deadline = Instant::now() + budget; + while kicked.is_ok() { + while let Some(completion) = self.receive() { + self.ready.push_back(completion); + } + if self.in_flight == self.ready.len() { + break; + } + if Instant::now() >= deadline { + // Too slow to wait for. Completions popped later say when that changes. + self.poll_pays = false; + break; + } + std::hint::spin_loop(); + } + self.received.polling.store(false, Ordering::SeqCst); + // What the worker finished before it could see the flag cleared went unsignalled. + while let Some(completion) = self.receive() { + self.ready.push_back(completion); + } + + kicked.map(|_| !self.ready.is_empty()) + } + /// Wait for every submitted request to complete. Their completions are left to be popped, /// unless `discard` is set. pub fn drain(&mut self, discard: bool) -> Result<(), ThreadedIoError> { @@ -770,6 +894,67 @@ mod tests { assert!(engine.pop().is_none()); } + #[test] + fn test_kick_and_poll() { + let mem = create_mem(); + let mut engine = new_engine(); + let push = |engine: &mut ThreadedFileEngine, count| { + for _ in 0..count { + engine + .push_write( + 0, + &mem, + GuestAddress(0), + FILE_LEN, + PendingRequest::default(), + ) + .unwrap(); + } + }; + + // With time to wait, everything in flight completes in the call, unsignalled. + engine.set_poll_budget(Some(Duration::from_secs(10))); + push(&mut engine, 4); + assert!(engine.kick_and_poll().unwrap()); + engine.completion_evt().read().unwrap_err(); + for _ in 0..4 { + assert_eq!(engine.pop().unwrap().result.unwrap(), FILE_LEN); + } + assert!(engine.pop().is_none()); + // Nothing in flight, nothing to wait for. + assert!(!engine.kick_and_poll().unwrap()); + + // With no time to wait, completions are signalled, and polling stops... + engine.set_poll_budget(Some(Duration::ZERO)); + push(&mut engine, 1); + assert!(!engine.kick_and_poll().unwrap()); + assert!(!engine.poll_pays); + wait_for_signal(&engine); + assert_eq!(engine.pop().unwrap().result.unwrap(), FILE_LEN); + assert!(!engine.poll_pays); + push(&mut engine, 1); + assert!(!engine.kick_and_poll().unwrap()); + wait_for_signal(&engine); + + // ...until a completion shows the backing file answers within the budget again. + engine.poll_budget = Some(Duration::from_secs(10)); + assert_eq!(engine.pop().unwrap().result.unwrap(), FILE_LEN); + assert!(engine.poll_pays); + push(&mut engine, 2); + assert!(engine.kick_and_poll().unwrap()); + engine.completion_evt().read().unwrap_err(); + assert_eq!(engine.in_flight, 2); + engine.drain(true).unwrap(); + assert_eq!(engine.in_flight, 0); + + // Turned off, a kick is only a kick. + engine.set_poll_budget(None); + push(&mut engine, 1); + assert!(!engine.kick_and_poll().unwrap()); + wait_for_signal(&engine); + assert_eq!(engine.pop().unwrap().result.unwrap(), FILE_LEN); + } + #[test] fn test_completions_are_signalled() { let mem = create_mem(); diff --git a/src/vmm/src/devices/virtio/block/virtio/threaded.rs b/src/vmm/src/devices/virtio/block/virtio/threaded.rs index 9550120d540..fa8cc9cb4da 100644 --- a/src/vmm/src/devices/virtio/block/virtio/threaded.rs +++ b/src/vmm/src/devices/virtio/block/virtio/threaded.rs @@ -226,6 +226,25 @@ impl VirtioBlock { } } + /// Hand the requests of a pass over the queue to the worker and, while the backing file is + /// fast, wait for them and complete them right away. + pub(crate) fn threaded_kick(&mut self) { + let FileEngine::Threaded(engine) = &mut self.disk.file_engine else { + return; + }; + // A throttled device resumes its queue from the completion event. + let completed = if self.is_io_engine_throttled { + engine.kick().map(|_| false) + } else { + engine.kick_and_poll() + }; + match completed { + Ok(true) => self.process_threaded_completion_queue(), + Ok(false) => {} + Err(err) => error!("BlockError submitting pending block requests: {:?}", err), + } + } + /// Add every request the worker finished to the used ring, and notify the guest. pub(crate) fn process_threaded_completion_queue(&mut self) { let FileEngine::Threaded(engine) = &mut self.disk.file_engine else { @@ -574,6 +593,36 @@ mod tests { )); } + #[test] + fn test_completed_in_the_queue_event() { + let mut block = default_block(FileEngineType::Threaded); + let FileEngine::Threaded(engine) = &mut block.disk.file_engine else { + unreachable!() + }; + engine.set_poll_budget(Some(std::time::Duration::from_secs(10))); + let mem = default_mem(); + let vq = VirtQueue::new(GuestAddress(0), &mem, 16); + set_queue(&mut block, 0, vq.create_queue()); + block.activate(mem.clone(), default_interrupt()).unwrap(); + + // The write is submitted, completed and notified in one queue event. + mem.write_obj::(123_456_789, GuestAddress(0x2000)) + .unwrap(); + let status_addr = set_segmented_request(&vq, VIRTIO_BLK_T_OUT, 0, &[(0x2000, 512)]); + simulate_queue_event(&mut block, Some(true)); + assert_eq!(vq.used.idx.get(), 1); + assert_eq!( + u32::from(mem.read_obj::(status_addr).unwrap()), + VIRTIO_BLK_S_OK + ); + // The completion event has nothing left to do. + let FileEngine::Threaded(engine) = &mut block.disk.file_engine else { + unreachable!() + }; + engine.completion_evt().read().unwrap_err(); + assert!(engine.pop().is_none()); + } + #[test] fn test_throttling() { let limit = u16::try_from(THREADED_IO_MAX_IN_FLIGHT).unwrap(); From 9fe071ca7c73cfb43b15b2c01e727b428ea6caee Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Tue, 29 Sep 2026 11:27:51 -0300 Subject: [PATCH 2/2] vmm: only wait for Threaded block completions with few in flight Waiting for completions in the queue event pays when the guest sends a request and does nothing until it completes, as it does for a write followed by fdatasync. With many requests in flight it does not: the completions are handled in bursts anyway, and the wait takes time from the event loop without adding throughput. Wait only when there are at most two requests to wait for. With a guest doing 64 KiB O_DIRECT writes at queue depth 32 on a tmpfs-backed drive, the event loop goes from 46-50% of a core back to 36-37%, which is what it takes without the wait, at the same 127-129k IOPS. Buffered writes with fdatasync keep the two wake-ups per commit. --- .../virtio/block/virtio/io/threaded_io.rs | 30 +++++++++++++++---- 1 file changed, 25 insertions(+), 5 deletions(-) diff --git a/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs b/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs index bfdc6a6f193..e6f3e89ec6a 100644 --- a/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs +++ b/src/vmm/src/devices/virtio/block/virtio/io/threaded_io.rs @@ -51,6 +51,11 @@ pub const THREADED_SEG_MAX: u32 = 32; /// loop serves nothing else, so it has to stay far below what any other device would notice. pub const POLL_BUDGET: Duration = Duration::from_micros(100); +/// How many requests may be in flight for the submitter to wait for them. Waiting pays when a +/// guest sends a request and does nothing until it completes. With more in flight, completions +/// are handled in bursts anyway, and the wait would only take time from the event loop. +pub const POLL_MAX_IN_FLIGHT: usize = 2; + /// A guest buffer: its address and length. pub type Segment = (GuestAddress, u32); @@ -582,13 +587,17 @@ impl ThreadedFileEngine { } /// Hand the requests pushed since the last kick to the worker, then wait for the completions - /// of everything in flight, for up to the poll budget, and only while the backing file has - /// been answering within it. Returns whether there are completions to pop. + /// of everything in flight, for up to the poll budget. Only while the backing file has been + /// answering within it, and with no more than [`POLL_MAX_IN_FLIGHT`] requests to wait for. + /// Returns whether there are completions to pop. /// /// The worker does not signal the completions it finishes meanwhile. pub fn kick_and_poll(&mut self) -> Result { + let waiting_for = self.in_flight - self.ready.len(); let budget = match self.poll_budget { - Some(budget) if self.poll_pays && self.in_flight > self.ready.len() => budget, + Some(budget) if self.poll_pays && (1..=POLL_MAX_IN_FLIGHT).contains(&waiting_for) => { + budget + } _ => { self.kick()?; return Ok(!self.ready.is_empty()); @@ -914,13 +923,24 @@ mod tests { // With time to wait, everything in flight completes in the call, unsignalled. engine.set_poll_budget(Some(Duration::from_secs(10))); - push(&mut engine, 4); + push(&mut engine, POLL_MAX_IN_FLIGHT); assert!(engine.kick_and_poll().unwrap()); engine.completion_evt().read().unwrap_err(); - for _ in 0..4 { + for _ in 0..POLL_MAX_IN_FLIGHT { assert_eq!(engine.pop().unwrap().result.unwrap(), FILE_LEN); } assert!(engine.pop().is_none()); + + // With more in flight than that, completions are signalled, and polling stays on. + push(&mut engine, POLL_MAX_IN_FLIGHT + 1); + assert!(!engine.kick_and_poll().unwrap()); + assert!(engine.poll_pays); + engine.drain(false).unwrap(); + wait_for_signal(&engine); + for _ in 0..=POLL_MAX_IN_FLIGHT { + assert_eq!(engine.pop().unwrap().result.unwrap(), FILE_LEN); + } + assert!(engine.poll_pays); // Nothing in flight, nothing to wait for. assert!(!engine.kick_and_poll().unwrap());