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..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 @@ -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,15 @@ 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); + +/// 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); @@ -86,6 +101,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 +274,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 +302,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 +323,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 +331,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 +369,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 +397,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 +418,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 +439,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 +473,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 +556,81 @@ 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. 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 && (1..=POLL_MAX_IN_FLIGHT).contains(&waiting_for) => { + 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 +903,78 @@ 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, POLL_MAX_IN_FLIGHT); + assert!(engine.kick_and_poll().unwrap()); + engine.completion_evt().read().unwrap_err(); + 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()); + + // 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();