From 2d3ac70dca24ecc5c495b3386aec6389e4a0926f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Wed, 23 Sep 2026 17:23:39 -0300 Subject: [PATCH 1/2] vmm: run Sync block I/O on a per-drive worker thread Every virtio device is serviced by the single VMM event loop thread. With io_engine=Sync, the block device did its reads, writes and fsyncs inline on that thread, and drained the whole avail ring before returning to epoll. While the host storage behind a drive is slow, nothing else on the event loop runs: virtio-net in particular stops being serviced, so guest network latency climbs to the disk's latency and throughput collapses. The drive rate limiter does not help, since it only paces requests; each one still blocks the loop for as long as the disk takes. Keep FileEngineType::Sync as the API and snapshot value, but back it with a worker thread per drive. The device submits requests to it and handles completions through an eventfd, exactly like the Async engine does: - The worker owns a clone of the backing file and performs requests one at a time, in submission order, so flushes still cover the writes submitted before them. Reads mark the guest pages they fill as dirty. - At most 128 requests are in flight; beyond that the engine reports itself throttled and the device resumes the queue on completion, as with io_uring. - drain()/drain_and_flush() wait for the worker through a barrier, so snapshotting and device teardown see every request finished. - A drive update sends the new file through the same queue: requests submitted before it use the old file, later ones the new one. The worker is started when the drive is created, before the VMM thread installs its seccomp filter, so it installs its own: a new, optional "block_io" thread category with the few syscalls it needs. Custom filter files without it fall back to the "vmm" filter, which is what this I/O ran under before. Measured with a guest writing to a drive whose host backing file adds 2 ms per I/O, while the host streams TCP into the guest and pings it: net under load RTT p99 before 68 MB/s 288 ms after 3125 MB/s 0.4 ms Guest disk throughput is unchanged (~440 iops). On a tmpfs-backed drive, guest fio (64k randwrite, QD32) goes from ~24k to ~107k iops, since the vCPU no longer waits for each request to complete when it kicks the queue. --- .../seccomp/aarch64-unknown-linux-musl.json | 163 ++++++++ resources/seccomp/unimplemented.json | 5 + .../seccomp/x86_64-unknown-linux-musl.json | 163 ++++++++ src/firecracker/src/main.rs | 8 + src/firecracker/src/seccomp.rs | 35 +- .../src/devices/virtio/block/virtio/device.rs | 68 ++-- .../virtio/block/virtio/event_handler.rs | 13 +- .../src/devices/virtio/block/virtio/io/mod.rs | 350 ++++++++++++------ .../devices/virtio/block/virtio/io/sync_io.rs | 314 ++++++++++++++-- .../devices/virtio/block/virtio/request.rs | 3 - .../devices/virtio/block/virtio/test_utils.rs | 21 +- src/vmm/src/seccomp.rs | 20 +- 12 files changed, 950 insertions(+), 213 deletions(-) diff --git a/resources/seccomp/aarch64-unknown-linux-musl.json b/resources/seccomp/aarch64-unknown-linux-musl.json index 26dd661e46b..0e105b0d9c9 100644 --- a/resources/seccomp/aarch64-unknown-linux-musl.json +++ b/resources/seccomp/aarch64-unknown-linux-musl.json @@ -1094,5 +1094,168 @@ ] } ] + }, + "block_io": { + "default_action": "trap", + "filter_action": "allow", + "filter": [ + { + "syscall": "exit" + }, + { + "syscall": "exit_group" + }, + { + "syscall": "read", + "comment": "Reads from the backing file" + }, + { + "syscall": "write", + "comment": "Writes to the backing file, and signals completions through an eventfd" + }, + { + "syscall": "fsync", + "comment": "Flush requests" + }, + { + "syscall": "close", + "comment": "Closes the previous backing file after a drive update" + }, + { + "syscall": "brk", + "comment": "Called for expanding the heap" + }, + { + "syscall": "clock_gettime", + "comment": "Used for metrics and logging, via the helpers in utils/src/time.rs. It's not called on some platforms, because of vdso optimisations." + }, + { + "syscall": "lseek", + "comment": "Positions the backing file before each read or write" + }, + { + "syscall": "mremap", + "comment": "Used for re-allocating large memory regions, for example vectors" + }, + { + "syscall": "munmap", + "comment": "Used for freeing memory" + }, + { + "syscall": "rt_sigprocmask", + "comment": "rt_sigprocmask is used by libc::abort during a panic to block and unblock signals" + }, + { + "syscall": "rt_sigreturn", + "comment": "rt_sigreturn is needed in case a fault does occur, so that the signal handler can return. Otherwise we get stuck in a fault loop." + }, + { + "syscall": "sigaltstack", + "comment": "sigaltstack is used by Rust stdlib to remove alternative signal stack during thread teardown." + }, + { + "syscall": "futex", + "comment": "Used for synchronization (during thread teardown when joining multiple vcpu threads at once)", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 0, + "comment": "FUTEX_WAIT" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization (during thread teardown)", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 1, + "comment": "FUTEX_WAKE" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 128, + "comment": "FUTEX_WAIT_PRIVATE" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 137, + "comment": "FUTEX_WAIT_BITSET_PRIVATE" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 129, + "comment": "FUTEX_WAKE_PRIVATE" + } + ] + }, + { + "syscall": "madvise", + "comment": "Used by the VirtIO balloon device and by musl for some customer workloads. It is also used by aws-lc during random number generation. They setup a memory page that mark with MADV_WIPEONFORK to be able to detect forks. They also call it with -1 to see if madvise is supported in certain platforms." + }, + { + "syscall": "mmap", + "comment": "Used by the allocator", + "args": [ + { + "index": 3, + "type": "dword", + "op": "eq", + "val": 34, + "comment": "libc::MAP_ANONYMOUS | libc::MAP_PRIVATE" + } + ] + }, + { + "syscall": "tkill", + "comment": "tkill is used by libc::abort during a panic to raise SIGABRT", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 6, + "comment": "SIGABRT" + } + ] + }, + { + "syscall": "sched_yield", + "comment": "Used by the rust standard library in std::sync::mpmc. Firecracker uses mpsc channels from this module for inter-thread communication" + }, + { + "syscall": "restart_syscall", + "comment": "automatically issued by the kernel when specific timing-related syscalls (e.g. nanosleep) get interrupted by SIGSTOP" + } + ] } } diff --git a/resources/seccomp/unimplemented.json b/resources/seccomp/unimplemented.json index a919df15519..f733949931f 100644 --- a/resources/seccomp/unimplemented.json +++ b/resources/seccomp/unimplemented.json @@ -13,5 +13,10 @@ "default_action": "allow", "filter_action": "trap", "filter": [] + }, + "block_io": { + "default_action": "allow", + "filter_action": "trap", + "filter": [] } } diff --git a/resources/seccomp/x86_64-unknown-linux-musl.json b/resources/seccomp/x86_64-unknown-linux-musl.json index dcd6753a4c5..c742f6e6f0c 100644 --- a/resources/seccomp/x86_64-unknown-linux-musl.json +++ b/resources/seccomp/x86_64-unknown-linux-musl.json @@ -1226,5 +1226,168 @@ ] } ] + }, + "block_io": { + "default_action": "trap", + "filter_action": "allow", + "filter": [ + { + "syscall": "exit" + }, + { + "syscall": "exit_group" + }, + { + "syscall": "read", + "comment": "Reads from the backing file" + }, + { + "syscall": "write", + "comment": "Writes to the backing file, and signals completions through an eventfd" + }, + { + "syscall": "fsync", + "comment": "Flush requests" + }, + { + "syscall": "close", + "comment": "Closes the previous backing file after a drive update" + }, + { + "syscall": "brk", + "comment": "Called for expanding the heap" + }, + { + "syscall": "clock_gettime", + "comment": "Used for metrics and logging, via the helpers in utils/src/time.rs. It's not called on some platforms, because of vdso optimisations." + }, + { + "syscall": "lseek", + "comment": "Positions the backing file before each read or write" + }, + { + "syscall": "mremap", + "comment": "Used for re-allocating large memory regions, for example vectors" + }, + { + "syscall": "munmap", + "comment": "Used for freeing memory" + }, + { + "syscall": "rt_sigprocmask", + "comment": "rt_sigprocmask is used by libc::abort during a panic to block and unblock signals" + }, + { + "syscall": "rt_sigreturn", + "comment": "rt_sigreturn is needed in case a fault does occur, so that the signal handler can return. Otherwise we get stuck in a fault loop." + }, + { + "syscall": "sigaltstack", + "comment": "sigaltstack is used by Rust stdlib to remove alternative signal stack during thread teardown." + }, + { + "syscall": "futex", + "comment": "Used for synchronization (during thread teardown when joining multiple vcpu threads at once)", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 0, + "comment": "FUTEX_WAIT" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization (during thread teardown)", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 1, + "comment": "FUTEX_WAKE" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 128, + "comment": "FUTEX_WAIT_PRIVATE" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 137, + "comment": "FUTEX_WAIT_BITSET_PRIVATE" + } + ] + }, + { + "syscall": "futex", + "comment": "Used for synchronization", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 129, + "comment": "FUTEX_WAKE_PRIVATE" + } + ] + }, + { + "syscall": "madvise", + "comment": "Used by the VirtIO balloon device and by musl for some customer workloads. It is also used by aws-lc during random number generation. They setup a memory page that mark with MADV_WIPEONFORK to be able to detect forks. They also call it with -1 to see if madvise is supported in certain platforms." + }, + { + "syscall": "mmap", + "comment": "Used by the allocator", + "args": [ + { + "index": 3, + "type": "dword", + "op": "eq", + "val": 34, + "comment": "libc::MAP_ANONYMOUS | libc::MAP_PRIVATE" + } + ] + }, + { + "syscall": "tkill", + "comment": "tkill is used by libc::abort during a panic to raise SIGABRT", + "args": [ + { + "index": 1, + "type": "dword", + "op": "eq", + "val": 6, + "comment": "SIGABRT" + } + ] + }, + { + "syscall": "sched_yield", + "comment": "Used by the rust standard library in std::sync::mpmc. Firecracker uses mpsc channels from this module for inter-thread communication" + }, + { + "syscall": "restart_syscall", + "comment": "automatically issued by the kernel when specific timing-related syscalls (e.g. nanosleep) get interrupted by SIGSTOP" + } + ] } } diff --git a/src/firecracker/src/main.rs b/src/firecracker/src/main.rs index 739214999a4..b240dfddccf 100644 --- a/src/firecracker/src/main.rs +++ b/src/firecracker/src/main.rs @@ -363,6 +363,14 @@ fn main_exec() -> Result<(), MainError> { .and_then(seccomp::get_filters) .map_err(MainError::SeccompFilter)?; + // Block device IO worker threads install their own filter; see `vmm::seccomp`. + if let Some(filter) = seccomp_filters + .get("block_io") + .or_else(|| seccomp_filters.get("vmm")) + { + vmm::seccomp::set_block_io_filter(filter.clone()); + } + let vmm_config_json = arguments .single_value("config-file") .map(fs::read_to_string) diff --git a/src/firecracker/src/seccomp.rs b/src/firecracker/src/seccomp.rs index 421220a7b5f..c0ffb774343 100644 --- a/src/firecracker/src/seccomp.rs +++ b/src/firecracker/src/seccomp.rs @@ -8,6 +8,9 @@ use std::path::Path; use vmm::seccomp::{BpfThreadMap, DeserializationError, deserialize_binary, get_empty_filters}; const THREAD_CATEGORIES: [&str; 3] = ["vmm", "api", "vcpu"]; +/// Thread categories a filter file may leave out. A missing "block_io" filter falls back to the +/// "vmm" one: block IO used to run on the VMM thread, so that filter already allows it. +const OPTIONAL_THREAD_CATEGORIES: [&str; 1] = ["block_io"]; /// Error retrieving seccomp filters. #[derive(Debug, thiserror::Error, displaydoc::Display)] @@ -78,9 +81,11 @@ fn get_custom_filters(reader: R) -> Result Result { - let (filters, invalid_filters): (BpfThreadMap, BpfThreadMap) = map - .into_iter() - .partition(|(k, _)| THREAD_CATEGORIES.contains(&k.as_str())); + let (filters, invalid_filters): (BpfThreadMap, BpfThreadMap) = + map.into_iter().partition(|(k, _)| { + THREAD_CATEGORIES.contains(&k.as_str()) + || OPTIONAL_THREAD_CATEGORIES.contains(&k.as_str()) + }); if !invalid_filters.is_empty() { // build the error message let mut thread_categories_string = @@ -117,16 +122,18 @@ mod tests { #[test] fn test_get_filters() { let mut filters = get_empty_filters(); - assert_eq!(filters.len(), 3); + assert_eq!(filters.len(), 4); assert!(filters.remove("vmm").is_some()); assert!(filters.remove("api").is_some()); assert!(filters.remove("vcpu").is_some()); + assert!(filters.remove("block_io").is_some()); let mut filters = get_empty_filters(); - assert_eq!(filters.len(), 3); + assert_eq!(filters.len(), 4); assert_eq!(filters.remove("vmm").unwrap().len(), 0); assert_eq!(filters.remove("api").unwrap().len(), 0); assert_eq!(filters.remove("vcpu").unwrap().len(), 0); + assert_eq!(filters.remove("block_io").unwrap().len(), 0); let file = TempFile::new().unwrap().into_file(); @@ -143,6 +150,15 @@ mod tests { assert_eq!(filter_thread_categories(map).unwrap().len(), 3); + // correct categories, including the optional ones + let mut map = BpfThreadMap::new(); + map.insert("vcpu".to_string(), Arc::new(vec![])); + map.insert("vmm".to_string(), Arc::new(vec![])); + map.insert("api".to_string(), Arc::new(vec![])); + map.insert("block_io".to_string(), Arc::new(vec![])); + + assert_eq!(filter_thread_categories(map).unwrap().len(), 4); + // invalid categories let mut map = BpfThreadMap::new(); map.insert("vcpu".to_string(), Arc::new(vec![])); @@ -168,6 +184,15 @@ mod tests { } } + #[test] + fn test_default_filters() { + // Debug builds compile an empty policy, so only release builds check the real one. + let filters = get_default_filters().unwrap(); + for category in THREAD_CATEGORIES.iter().chain(&OPTIONAL_THREAD_CATEGORIES) { + assert!(filters.contains_key(*category), "missing {category}"); + } + } + #[test] fn test_seccomp_config() { assert!(matches!( diff --git a/src/vmm/src/devices/virtio/block/virtio/device.rs b/src/vmm/src/devices/virtio/block/virtio/device.rs index ecdd8ee4f6d..b0e773e0a73 100644 --- a/src/vmm/src/devices/virtio/block/virtio/device.rs +++ b/src/vmm/src/devices/virtio/block/virtio/device.rs @@ -19,7 +19,6 @@ use serde::{Deserialize, Serialize}; use vm_memory::ByteValued; use vmm_sys_util::eventfd::EventFd; -use super::io::async_io; use super::request::*; use super::{BLOCK_QUEUE_SIZES, SECTOR_SHIFT, SECTOR_SIZE, VirtioBlockError, io as block_io}; use crate::devices::virtio::ActivateError; @@ -268,18 +267,6 @@ pub struct VirtioBlock { pub metrics: Arc, } -macro_rules! unwrap_async_file_engine_or_return { - ($file_engine: expr) => { - match $file_engine { - FileEngine::Async(engine) => engine, - FileEngine::Sync(_) => { - error!("The block device doesn't use an async IO engine"); - return; - } - } - }; -} - impl VirtioBlock { /// Create a new virtio block device that operates on the given file. /// @@ -468,34 +455,22 @@ impl VirtioBlock { Ok(()) } + /// Hand the requests the IO engine finished back to the guest. fn process_async_completion_queue(&mut self) { - let engine = unwrap_async_file_engine_or_return!(&mut self.disk.file_engine); - // This is safe since we checked in the event handler that the device is activated. let active_state = self.device_state.active_state().unwrap(); let queue = &mut self.queues[0]; loop { - match engine.pop(&active_state.mem) { + match self.disk.file_engine.pop(&active_state.mem) { Err(error) => { - error!("Failed to read completed io_uring entry: {:?}", error); + error!("Failed to read completed block request: {:?}", error); break; } Ok(None) => break, - Ok(Some(cqe)) => { - let res = cqe.result(); - let user_data = cqe.user_data(); - - let (pending, res) = match res { - Ok(count) => (user_data, Ok(count)), - Err(error) => ( - user_data, - Err(IoErr::FileEngine(block_io::BlockIoError::Async( - async_io::AsyncIoError::IO(error), - ))), - ), - }; - let finished = pending.finish(&active_state.mem, res, &self.metrics); + Ok(Some(completion)) => { + let res = completion.result.map_err(IoErr::FileEngine); + let finished = completion.req.finish(&active_state.mem, res, &self.metrics); queue .add_used(finished.desc_idx, finished.num_bytes_to_mem) .unwrap_or_else(|err| { @@ -520,9 +495,7 @@ impl VirtioBlock { } pub fn process_async_completion_event(&mut self) { - let engine = unwrap_async_file_engine_or_return!(&mut self.disk.file_engine); - - if let Err(err) = engine.completion_evt().read() { + if let Err(err) = self.disk.file_engine.completion_evt().read() { error!("Failed to get async completion event: {:?}", err); } else { self.process_async_completion_queue(); @@ -576,9 +549,7 @@ impl VirtioBlock { } self.drain_and_flush(false); - if let FileEngine::Async(ref _engine) = self.disk.file_engine { - self.process_async_completion_queue(); - } + self.process_async_completion_queue(); } } @@ -701,6 +672,7 @@ mod tests { use super::*; use crate::check_metric_after_block; use crate::devices::virtio::block::virtio::IO_URING_NUM_ENTRIES; + use crate::devices::virtio::block::virtio::io::sync_io::SYNC_IO_MAX_IN_FLIGHT; use crate::devices::virtio::block::virtio::test_utils::{ default_block, read_blk_req_descriptors, set_queue, set_rate_limiter, simulate_async_completion_event, simulate_queue_and_async_completion_events, @@ -1594,25 +1566,29 @@ mod tests { #[test] fn test_io_engine_throttling() { - // FullSQueue BlockError - { - let mut block = default_block(FileEngineType::Async); + // FullSQueue BlockError, or the sync engine's in-flight limit. + let sync_limit = u16::try_from(SYNC_IO_MAX_IN_FLIGHT).unwrap(); + for (engine, limit) in [ + (FileEngineType::Sync, sync_limit), + (FileEngineType::Async, IO_URING_NUM_ENTRIES), + ] { + let mut block = default_block(engine); let mem = default_mem(); let interrupt = default_interrupt(); - let vq = VirtQueue::new(GuestAddress(0), &mem, IO_URING_NUM_ENTRIES * 4); + let vq = VirtQueue::new(GuestAddress(0), &mem, limit * 4); block.queues[0] = vq.create_queue(); block.activate(mem.clone(), interrupt).unwrap(); // Run scenario that doesn't trigger FullSq BlockError: Add sq_size flush requests. - add_flush_requests_batch(&mut block, &vq, IO_URING_NUM_ENTRIES); + add_flush_requests_batch(&mut block, &vq, limit); simulate_queue_event(&mut block, Some(false)); assert!(!block.is_io_engine_throttled); simulate_async_completion_event(&mut block, true); - check_flush_requests_batch(IO_URING_NUM_ENTRIES, &vq); + check_flush_requests_batch(limit, &vq); // Run scenario that triggers FullSqError : Add sq_size + 10 flush requests. - add_flush_requests_batch(&mut block, &vq, IO_URING_NUM_ENTRIES + 10); + add_flush_requests_batch(&mut block, &vq, limit + 10); simulate_queue_event(&mut block, Some(false)); assert!(block.is_io_engine_throttled); // When the async_completion_event is triggered: @@ -1621,12 +1597,12 @@ mod tests { // 3. process_queue() should be called again. simulate_async_completion_event(&mut block, true); assert!(!block.is_io_engine_throttled); - check_flush_requests_batch(IO_URING_NUM_ENTRIES, &vq); + check_flush_requests_batch(limit, &vq); // check that process_queue() was called again resulting in the processing of the // remaining 10 ops. simulate_async_completion_event(&mut block, true); assert!(!block.is_io_engine_throttled); - check_flush_requests_batch(IO_URING_NUM_ENTRIES + 10, &vq); + check_flush_requests_batch(limit + 10, &vq); } // FullCQueue BlockError diff --git a/src/vmm/src/devices/virtio/block/virtio/event_handler.rs b/src/vmm/src/devices/virtio/block/virtio/event_handler.rs index 9f02862f814..6feee595d08 100644 --- a/src/vmm/src/devices/virtio/block/virtio/event_handler.rs +++ b/src/vmm/src/devices/virtio/block/virtio/event_handler.rs @@ -3,7 +3,6 @@ use event_manager::{EventOps, Events, MutEventSubscriber}; use vmm_sys_util::epoll::EventSet; -use super::io::FileEngine; use crate::devices::virtio::block::virtio::device::VirtioBlock; use crate::devices::virtio::device::VirtioDevice; use crate::logger::{error, warn}; @@ -29,13 +28,11 @@ impl VirtioBlock { )) { error!("Failed to register ratelimiter event: {}", err); } - if let FileEngine::Async(ref engine) = self.disk.file_engine - && let Err(err) = ops.add(Events::with_data( - engine.completion_evt(), - Self::PROCESS_ASYNC_COMPLETION, - EventSet::IN, - )) - { + if let Err(err) = ops.add(Events::with_data( + self.disk.file_engine.completion_evt(), + Self::PROCESS_ASYNC_COMPLETION, + EventSet::IN, + )) { error!("Failed to register IO engine completion event: {}", err); } } diff --git a/src/vmm/src/devices/virtio/block/virtio/io/mod.rs b/src/vmm/src/devices/virtio/block/virtio/io/mod.rs index b7aa8061d76..5367217daae 100644 --- a/src/vmm/src/devices/virtio/block/virtio/io/mod.rs +++ b/src/vmm/src/devices/virtio/block/virtio/io/mod.rs @@ -12,17 +12,11 @@ pub use self::sync_io::{SyncFileEngine, SyncIoError}; use crate::devices::virtio::block::virtio::PendingRequest; use crate::devices::virtio::block::virtio::device::FileEngineType; use crate::vstate::memory::{GuestAddress, GuestMemoryMmap}; - -#[derive(Debug)] -pub struct RequestOk { - pub req: PendingRequest, - pub count: u32, -} +use vmm_sys_util::eventfd::EventFd; #[derive(Debug)] pub enum FileEngineOk { Submitted, - Executed(RequestOk), } #[derive(Debug, thiserror::Error, displaydoc::Display)] @@ -37,6 +31,7 @@ impl BlockIoError { pub fn is_throttling_err(&self) -> bool { match self { BlockIoError::Async(AsyncIoError::IoUring(err)) => err.is_throttling_err(), + BlockIoError::Sync(SyncIoError::QueueFull) => true, _ => false, } } @@ -48,6 +43,22 @@ pub struct RequestError { pub error: E, } +impl RequestError { + fn map(self, f: impl FnOnce(E) -> F) -> RequestError { + RequestError { + req: self.req, + error: f(self.error), + } + } +} + +/// A request the engine finished, with its outcome. +#[derive(Debug)] +pub struct Completion { + pub req: PendingRequest, + pub result: Result, +} + #[allow(clippy::large_enum_variant)] #[derive(Debug)] pub enum FileEngine { @@ -62,14 +73,16 @@ impl FileEngine { FileEngineType::Async => Ok(FileEngine::Async( AsyncFileEngine::from_file(file).map_err(BlockIoError::Async)?, )), - FileEngineType::Sync => Ok(FileEngine::Sync(SyncFileEngine::from_file(file))), + FileEngineType::Sync => Ok(FileEngine::Sync( + SyncFileEngine::from_file(file).map_err(BlockIoError::Sync)?, + )), } } pub fn update_file_path(&mut self, file: File) -> Result<(), BlockIoError> { match self { FileEngine::Async(engine) => engine.update_file(file).map_err(BlockIoError::Async)?, - FileEngine::Sync(engine) => engine.update_file(file), + FileEngine::Sync(engine) => engine.update_file(file).map_err(BlockIoError::Sync)?, }; Ok(()) @@ -83,6 +96,14 @@ impl FileEngine { } } + /// The eventfd signalled when submitted requests complete. + pub fn completion_evt(&self) -> &EventFd { + match self { + FileEngine::Async(engine) => engine.completion_evt(), + FileEngine::Sync(engine) => engine.completion_evt(), + } + } + pub fn read( &mut self, offset: u64, @@ -92,21 +113,14 @@ impl FileEngine { req: PendingRequest, ) -> Result> { match self { - FileEngine::Async(engine) => match engine.push_read(offset, mem, addr, count, req) { - Ok(_) => Ok(FileEngineOk::Submitted), - Err(err) => Err(RequestError { - req: err.req, - error: BlockIoError::Async(err.error), - }), - }, - FileEngine::Sync(engine) => match engine.read(offset, mem, addr, count) { - Ok(count) => Ok(FileEngineOk::Executed(RequestOk { req, count })), - Err(err) => Err(RequestError { - req, - error: BlockIoError::Sync(err), - }), - }, + FileEngine::Async(engine) => engine + .push_read(offset, mem, addr, count, req) + .map_err(|err| err.map(BlockIoError::Async)), + FileEngine::Sync(engine) => engine + .push_read(offset, mem, addr, count, req) + .map_err(|err| err.map(BlockIoError::Sync)), } + .map(|_| FileEngineOk::Submitted) } pub fn write( @@ -118,21 +132,14 @@ impl FileEngine { req: PendingRequest, ) -> Result> { match self { - FileEngine::Async(engine) => match engine.push_write(offset, mem, addr, count, req) { - Ok(_) => Ok(FileEngineOk::Submitted), - Err(err) => Err(RequestError { - req: err.req, - error: BlockIoError::Async(err.error), - }), - }, - FileEngine::Sync(engine) => match engine.write(offset, mem, addr, count) { - Ok(count) => Ok(FileEngineOk::Executed(RequestOk { req, count })), - Err(err) => Err(RequestError { - req, - error: BlockIoError::Sync(err), - }), - }, + FileEngine::Async(engine) => engine + .push_write(offset, mem, addr, count, req) + .map_err(|err| err.map(BlockIoError::Async)), + FileEngine::Sync(engine) => engine + .push_write(offset, mem, addr, count, req) + .map_err(|err| err.map(BlockIoError::Sync)), } + .map(|_| FileEngineOk::Submitted) } pub fn flush( @@ -140,27 +147,42 @@ impl FileEngine { req: PendingRequest, ) -> Result> { match self { - FileEngine::Async(engine) => match engine.push_flush(req) { - Ok(_) => Ok(FileEngineOk::Submitted), - Err(err) => Err(RequestError { - req: err.req, - error: BlockIoError::Async(err.error), - }), - }, - FileEngine::Sync(engine) => match engine.flush() { - Ok(_) => Ok(FileEngineOk::Executed(RequestOk { req, count: 0 })), - Err(err) => Err(RequestError { - req, - error: BlockIoError::Sync(err), - }), - }, + FileEngine::Async(engine) => engine + .push_flush(req) + .map_err(|err| err.map(BlockIoError::Async)), + FileEngine::Sync(engine) => engine + .push_flush(req) + .map_err(|err| err.map(BlockIoError::Sync)), + } + .map(|_| FileEngineOk::Submitted) + } + + /// Pop a finished request, if there is one. + pub fn pop(&mut self, mem: &GuestMemoryMmap) -> Result, BlockIoError> { + match self { + FileEngine::Async(engine) => { + let cqe = engine.pop(mem).map_err(BlockIoError::Async)?; + Ok(cqe.map(|cqe| { + let result = cqe + .result() + .map_err(|err| BlockIoError::Async(AsyncIoError::IO(err))); + Completion { + req: cqe.user_data(), + result, + } + })) + } + FileEngine::Sync(engine) => Ok(engine.pop().map(|completion| Completion { + req: completion.req, + result: completion.result.map_err(BlockIoError::Sync), + })), } } pub fn drain(&mut self, discard: bool) -> Result<(), BlockIoError> { match self { FileEngine::Async(engine) => engine.drain(discard).map_err(BlockIoError::Async), - FileEngine::Sync(_engine) => Ok(()), + FileEngine::Sync(engine) => engine.drain(discard).map_err(BlockIoError::Sync), } } @@ -169,7 +191,7 @@ impl FileEngine { FileEngine::Async(engine) => { engine.drain_and_flush(discard).map_err(BlockIoError::Async) } - FileEngine::Sync(engine) => engine.flush().map_err(BlockIoError::Sync), + FileEngine::Sync(engine) => engine.drain_and_flush(discard).map_err(BlockIoError::Sync), } } } @@ -178,6 +200,7 @@ impl FileEngine { pub mod tests { #![allow(clippy::undocumented_unsafe_blocks)] use std::os::unix::ffi::OsStrExt; + use std::os::unix::fs::MetadataExt; use vm_memory::GuestMemoryRegion; use vmm_sys_util::tempfile::TempFile; @@ -193,32 +216,16 @@ pub mod tests { // 2 pages of memory should be enough to test read/write ops and also dirty tracking. const MEM_LEN: usize = 8192; - macro_rules! assert_sync_execution { - ($expression:expr, $count:expr) => { - match $expression { - Ok(FileEngineOk::Executed(RequestOk { req: _, count })) => { - assert_eq!(count, $count) - } - other => panic!( - "Expected: Ok(FileEngineOk::Executed(UserDataOk {{ user_data: _, count: {} \ - }})), got: {:?}", - $count, other - ), - } - }; - } - macro_rules! assert_queued { ($expression:expr) => { assert!(matches!($expression, Ok(FileEngineOk::Submitted))) }; } - fn assert_async_execution(mem: &GuestMemoryMmap, engine: &mut FileEngine, count: u32) { - if let FileEngine::Async(engine) = engine { - engine.drain(false).unwrap(); - assert_eq!(engine.pop(mem).unwrap().unwrap().result().unwrap(), count); - } + fn assert_execution(mem: &GuestMemoryMmap, engine: &mut FileEngine, count: u32) { + engine.drain(false).unwrap(); + let completion = engine.pop(mem).unwrap().unwrap(); + assert_eq!(completion.result.unwrap(), count); } fn create_mem() -> GuestMemoryMmap { @@ -265,16 +272,12 @@ pub mod tests { let partial_len = 50; let addr = GuestAddress(MEM_LEN as u64 - u64::from(partial_len)); mem.write(&data, addr).unwrap(); - assert_sync_execution!( - engine.write(0, &mem, addr, partial_len, PendingRequest::default()), - partial_len - ); + assert_queued!(engine.write(0, &mem, addr, partial_len, PendingRequest::default())); + assert_execution(&mem, &mut engine, partial_len); // Partial read let mem = create_mem(); - assert_sync_execution!( - engine.read(0, &mem, addr, partial_len, PendingRequest::default()), - partial_len - ); + assert_queued!(engine.read(0, &mem, addr, partial_len, PendingRequest::default())); + assert_execution(&mem, &mut engine, partial_len); // Check data let mut buf = vec![0u8; partial_len as usize]; mem.read_slice(&mut buf, addr).unwrap(); @@ -285,56 +288,173 @@ pub mod tests { let partial_len = 50; let addr = GuestAddress(0); mem.write(&data, addr).unwrap(); - assert_sync_execution!( - engine.write(offset, &mem, addr, partial_len, PendingRequest::default()), - partial_len - ); + assert_queued!(engine.write(offset, &mem, addr, partial_len, PendingRequest::default())); + assert_execution(&mem, &mut engine, partial_len); // Offset read let mem = create_mem(); - assert_sync_execution!( - engine.read(offset, &mem, addr, partial_len, PendingRequest::default()), - partial_len - ); + assert_queued!(engine.read(offset, &mem, addr, partial_len, PendingRequest::default())); + assert_execution(&mem, &mut engine, partial_len); // Check data let mut buf = vec![0u8; partial_len as usize]; mem.read_slice(&mut buf, addr).unwrap(); assert_eq!(buf, data[..partial_len as usize]); + // check dirty mem + check_dirty_mem(&mem, addr, partial_len); + check_clean_mem(&mem, GuestAddress(4096), 4096); // Full write mem.write(&data, GuestAddress(0)).unwrap(); - assert_sync_execution!( - engine.write( - 0, - &mem, - GuestAddress(0), - FILE_LEN, - PendingRequest::default() - ), - FILE_LEN - ); + assert_queued!(engine.write( + 0, + &mem, + GuestAddress(0), + FILE_LEN, + PendingRequest::default() + )); + assert_execution(&mem, &mut engine, FILE_LEN); // Full read let mem = create_mem(); - assert_sync_execution!( - engine.read( - 0, - &mem, - GuestAddress(0), - FILE_LEN, - PendingRequest::default() - ), - FILE_LEN - ); + assert_queued!(engine.read( + 0, + &mem, + GuestAddress(0), + FILE_LEN, + PendingRequest::default() + )); + assert_execution(&mem, &mut engine, FILE_LEN); // Check data let mut buf = vec![0u8; FILE_LEN as usize]; mem.read_slice(&mut buf, GuestAddress(0)).unwrap(); assert_eq!(buf, data.as_slice()); + // check dirty mem + check_dirty_mem(&mem, GuestAddress(0), FILE_LEN); + check_clean_mem(&mem, GuestAddress(4096), 4096); + + // Out of bounds guest memory fails the request, not the engine. + assert_queued!(engine.read( + 0, + &mem, + GuestAddress(MEM_LEN as u64), + FILE_LEN, + PendingRequest::default() + )); + engine.drain(false).unwrap(); + let completion = engine.pop(&mem).unwrap().unwrap(); + assert!(matches!( + completion.result, + Err(BlockIoError::Sync(SyncIoError::Transfer(_))) + )); // Check other ops - engine.flush(PendingRequest::default()).unwrap(); + assert_queued!(engine.flush(PendingRequest::default())); + assert_execution(&mem, &mut engine, 0); + engine.drain(true).unwrap(); engine.drain_and_flush(true).unwrap(); } + #[test] + fn test_sync_runs_off_thread() { + let mem = create_mem(); + let file = TempFile::new().unwrap().into_file(); + let mut engine = FileEngine::from_file(file, FileEngineType::Sync).unwrap(); + + // Submitting returns before the IO is done: completions only show up through the + // completion eventfd, once the worker has finished them. + for _ in 0..10 { + assert_queued!(engine.write( + 0, + &mem, + GuestAddress(0), + FILE_LEN, + PendingRequest::default() + )); + } + let mut completed = 0; + while completed < 10 { + // Blocks until the worker signals; the eventfd is non-blocking, so poll it. + match engine.completion_evt().read() { + Ok(_) => { + while let Some(completion) = engine.pop(&mem).unwrap() { + assert_eq!(completion.result.unwrap(), FILE_LEN); + completed += 1; + } + } + Err(err) => { + assert_eq!(err.kind(), std::io::ErrorKind::WouldBlock); + std::thread::yield_now(); + } + } + } + assert!(engine.pop(&mem).unwrap().is_none()); + } + + #[test] + fn test_sync_throttling() { + let mem = create_mem(); + let file = TempFile::new().unwrap().into_file(); + let mut engine = FileEngine::from_file(file, FileEngineType::Sync).unwrap(); + + for _ in 0..sync_io::SYNC_IO_MAX_IN_FLIGHT { + assert_queued!(engine.flush(PendingRequest::default())); + } + // Completed but not yet popped requests still count as in flight. + engine.drain(false).unwrap(); + let err = engine.flush(PendingRequest::default()).unwrap_err(); + assert!(err.error.is_throttling_err()); + + // Popping a completion makes room for one more request. + engine.pop(&mem).unwrap().unwrap(); + assert_queued!(engine.flush(PendingRequest::default())); + let err = engine.flush(PendingRequest::default()).unwrap_err(); + assert!(err.error.is_throttling_err()); + + // Discarding all completions makes room for all of them. + engine.drain(true).unwrap(); + assert!(engine.pop(&mem).unwrap().is_none()); + for _ in 0..sync_io::SYNC_IO_MAX_IN_FLIGHT { + assert_queued!(engine.flush(PendingRequest::default())); + } + engine.drain(true).unwrap(); + } + + #[test] + fn test_sync_update_file() { + let mem = create_mem(); + let old = TempFile::new().unwrap(); + let new = TempFile::new().unwrap(); + let mut engine = + FileEngine::from_file(old.as_file().try_clone().unwrap(), FileEngineType::Sync) + .unwrap(); + + let data = vmm_sys_util::rand::rand_alphanumerics(FILE_LEN as usize) + .as_bytes() + .to_vec(); + mem.write(&data, GuestAddress(0)).unwrap(); + + // Requests submitted before the update go to the old file, the ones after to the new one, + // without waiting for the first to complete. + assert_queued!(engine.write( + 0, + &mem, + GuestAddress(0), + FILE_LEN, + PendingRequest::default() + )); + engine + .update_file_path(new.as_file().try_clone().unwrap()) + .unwrap(); + assert_queued!(engine.write(0, &mem, GuestAddress(0), 10, PendingRequest::default())); + engine.drain(true).unwrap(); + + assert_eq!(old.as_file().metadata().unwrap().len(), u64::from(FILE_LEN)); + assert_eq!(new.as_file().metadata().unwrap().len(), 10); + assert_eq!( + engine.file().metadata().unwrap().ino(), + new.as_file().metadata().unwrap().ino() + ); + } + #[test] fn test_async() { // Create backing file. @@ -355,11 +475,11 @@ pub mod tests { let addr = GuestAddress(0); mem.write(&data, addr).unwrap(); assert_queued!(engine.write(offset, &mem, addr, partial_len, PendingRequest::default())); - assert_async_execution(&mem, &mut engine, partial_len); + assert_execution(&mem, &mut engine, partial_len); // Offset read let mem = create_mem(); assert_queued!(engine.read(offset, &mem, addr, partial_len, PendingRequest::default())); - assert_async_execution(&mem, &mut engine, partial_len); + assert_execution(&mem, &mut engine, partial_len); // Check data let mut buf = vec![0u8; partial_len as usize]; mem.read_slice(&mut buf, addr).unwrap(); @@ -371,12 +491,12 @@ pub mod tests { // Full write mem.write(&data, GuestAddress(0)).unwrap(); assert_queued!(engine.write(0, &mem, addr, FILE_LEN, PendingRequest::default())); - assert_async_execution(&mem, &mut engine, FILE_LEN); + assert_execution(&mem, &mut engine, FILE_LEN); // Full read let mem = create_mem(); assert_queued!(engine.read(0, &mem, addr, FILE_LEN, PendingRequest::default())); - assert_async_execution(&mem, &mut engine, FILE_LEN); + assert_execution(&mem, &mut engine, FILE_LEN); // Check data let mut buf = vec![0u8; FILE_LEN as usize]; mem.read_slice(&mut buf, GuestAddress(0)).unwrap(); @@ -387,7 +507,7 @@ pub mod tests { // Check other ops assert_queued!(engine.flush(PendingRequest::default())); - assert_async_execution(&mem, &mut engine, 0); + assert_execution(&mem, &mut engine, 0); engine.drain(true).unwrap(); engine.drain_and_flush(true).unwrap(); diff --git a/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs b/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs index eec3b3d8b8d..5495930a586 100644 --- a/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs +++ b/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs @@ -1,12 +1,30 @@ // Copyright 2021 Amazon.com, Inc. or its affiliates. All Rights Reserved. // SPDX-License-Identifier: Apache-2.0 +//! Sync file engine. +//! +//! The I/O itself is done with blocking system calls, but not on the thread that submits it: each +//! engine owns a worker thread that performs the requests one at a time, in submission order, and +//! reports each completion through an eventfd, the same way the async engine does. This keeps a +//! slow backing file from stalling the event loop, and with it every other device serviced there. + use std::fs::File; use std::io::{Seek, SeekFrom, Write}; +use std::sync::mpsc; +use std::thread; use vm_memory::{GuestMemoryError, ReadVolatile, WriteVolatile}; +use vmm_sys_util::eventfd::EventFd; + +use crate::devices::virtio::block::virtio::PendingRequest; +use crate::devices::virtio::block::virtio::io::RequestError; +use crate::logger::error; +use crate::vstate::memory::{GuestAddress, GuestMemory, GuestMemoryExtension, GuestMemoryMmap}; -use crate::vstate::memory::{GuestAddress, GuestMemory, GuestMemoryMmap}; +/// Maximum number of requests submitted to the worker and not yet popped. The engine reports +/// itself as throttled beyond this, and the device resumes processing its queue once completions +/// come back. +pub const SYNC_IO_MAX_IN_FLIGHT: usize = 128; #[derive(Debug, thiserror::Error, displaydoc::Display)] pub enum SyncIoError { @@ -18,32 +36,63 @@ pub enum SyncIoError { SyncAll(std::io::Error), /// Transfer: {0} Transfer(GuestMemoryError), + /// EventFd: {0} + EventFd(std::io::Error), + /// Cloning the backing file: {0} + FileClone(std::io::Error), + /// Spawning the IO worker thread: {0} + Spawn(std::io::Error), + /// Too many requests in flight + QueueFull, + /// The IO worker thread is gone + WorkerGone, } +/// A finished request, as reported by the worker thread. #[derive(Debug)] -pub struct SyncFileEngine { - file: File, +pub struct SyncCompletion { + pub req: PendingRequest, + pub result: Result, } -// SAFETY: `File` is send and ultimately a POD. -unsafe impl Send for SyncFileEngine {} - -impl SyncFileEngine { - pub fn from_file(file: File) -> SyncFileEngine { - SyncFileEngine { file } - } +#[derive(Debug)] +enum Io { + Read { + offset: u64, + mem: GuestMemoryMmap, + addr: GuestAddress, + count: u32, + }, + Write { + offset: u64, + mem: GuestMemoryMmap, + addr: GuestAddress, + count: u32, + }, + Flush, +} - #[cfg(test)] - pub fn file(&self) -> &File { - &self.file - } +#[derive(Debug)] +enum Op { + Io { + io: Io, + req: PendingRequest, + }, + /// Switch to a new backing file. Requests submitted before it use the old one. + UpdateFile(File), + /// Acknowledged once every request submitted before it has completed. + Barrier(mpsc::SyncSender<()>), + Exit, +} - /// Update the backing file of the engine - pub fn update_file(&mut self, file: File) { - self.file = file - } +/// Blocking I/O on the backing file. Only ever used from the worker thread. +#[derive(Debug)] +struct BlockingFile { + file: File, +} - pub fn read( +impl BlockingFile { + fn read( &mut self, offset: u64, mem: &GuestMemoryMmap, @@ -59,7 +108,7 @@ impl SyncFileEngine { Ok(count) } - pub fn write( + fn write( &mut self, offset: u64, mem: &GuestMemoryMmap, @@ -75,10 +124,233 @@ impl SyncFileEngine { Ok(count) } - pub fn flush(&mut self) -> Result<(), SyncIoError> { + fn flush(&mut self) -> Result<(), SyncIoError> { // flush() first to force any cached data out of rust buffers. self.file.flush().map_err(SyncIoError::Flush)?; // Sync data out to physical media on host. self.file.sync_all().map_err(SyncIoError::SyncAll) } + + fn execute(&mut self, io: Io) -> Result { + match io { + Io::Read { + offset, + mem, + addr, + count, + } => { + let count = self.read(offset, &mem, addr, count)?; + // The guest memory was written from this thread, so account for it in the dirty + // bitmap before the device gets to see the completion. + mem.mark_dirty(addr, count as usize); + Ok(count) + } + Io::Write { + offset, + mem, + addr, + count, + } => self.write(offset, &mem, addr, count), + Io::Flush => self.flush().map(|_| 0), + } + } +} + +fn run_worker( + mut file: BlockingFile, + ops: mpsc::Receiver, + completions: mpsc::Sender, + completion_evt: EventFd, +) { + if let Some(filter) = crate::seccomp::block_io_filter() + && let Err(err) = crate::seccomp::apply_filter(filter) + { + panic!("Failed to set the requested seccomp filters on the block IO worker: {err}"); + } + + for op in ops { + let completion = match op { + Op::Io { io, req } => SyncCompletion { + req, + result: file.execute(io), + }, + Op::UpdateFile(new_file) => { + file.file = new_file; + continue; + } + Op::Barrier(ack) => { + // The submitter may have given up waiting; nothing to do about it here. + let _ = ack.send(()); + continue; + } + Op::Exit => break, + }; + + if completions.send(completion).is_err() { + break; + } + if let Err(err) = completion_evt.write(1) { + error!("Failed to signal block IO completion: {:?}", err); + } + } +} + +/// Front end of the sync engine, used from the thread that owns the device. +#[derive(Debug)] +pub struct SyncFileEngine { + file: File, + ops: mpsc::Sender, + completions: mpsc::Receiver, + completion_evt: EventFd, + in_flight: usize, + worker: Option>, +} + +impl SyncFileEngine { + pub fn from_file(file: File) -> Result { + let completion_evt = EventFd::new(libc::EFD_NONBLOCK).map_err(SyncIoError::EventFd)?; + let worker_evt = completion_evt.try_clone().map_err(SyncIoError::EventFd)?; + let worker_file = BlockingFile { + file: file.try_clone().map_err(SyncIoError::FileClone)?, + }; + let (ops, worker_ops) = mpsc::channel(); + let (worker_completions, completions) = mpsc::channel(); + + let worker = thread::Builder::new() + .name("fc_blk_io".to_string()) + .spawn(move || run_worker(worker_file, worker_ops, worker_completions, worker_evt)) + .map_err(SyncIoError::Spawn)?; + + Ok(SyncFileEngine { + file, + ops, + completions, + completion_evt, + in_flight: 0, + worker: Some(worker), + }) + } + + #[cfg(test)] + pub fn file(&self) -> &File { + &self.file + } + + /// Update the backing file of the engine + pub fn update_file(&mut self, file: File) -> Result<(), SyncIoError> { + let worker_file = file.try_clone().map_err(SyncIoError::FileClone)?; + self.ops + .send(Op::UpdateFile(worker_file)) + .map_err(|_| SyncIoError::WorkerGone)?; + self.file = file; + Ok(()) + } + + pub fn completion_evt(&self) -> &EventFd { + &self.completion_evt + } + + fn push(&mut self, io: Io, req: PendingRequest) -> Result<(), RequestError> { + if self.in_flight >= SYNC_IO_MAX_IN_FLIGHT { + return Err(RequestError { + req, + error: SyncIoError::QueueFull, + }); + } + + match self.ops.send(Op::Io { io, req }) { + Ok(()) => { + self.in_flight += 1; + Ok(()) + } + Err(mpsc::SendError(Op::Io { req, .. })) => Err(RequestError { + req, + error: SyncIoError::WorkerGone, + }), + Err(_) => unreachable!("sent an IO op"), + } + } + + pub fn push_read( + &mut self, + offset: u64, + mem: &GuestMemoryMmap, + addr: GuestAddress, + count: u32, + req: PendingRequest, + ) -> Result<(), RequestError> { + let io = Io::Read { + offset, + mem: mem.clone(), + addr, + count, + }; + self.push(io, req) + } + + pub fn push_write( + &mut self, + offset: u64, + mem: &GuestMemoryMmap, + addr: GuestAddress, + count: u32, + req: PendingRequest, + ) -> Result<(), RequestError> { + let io = Io::Write { + offset, + mem: mem.clone(), + addr, + count, + }; + self.push(io, req) + } + + pub fn push_flush(&mut self, req: PendingRequest) -> Result<(), RequestError> { + self.push(Io::Flush, req) + } + + /// Pop a finished request, if there is one. + pub fn pop(&mut self) -> Option { + let completion = self.completions.try_recv().ok()?; + self.in_flight -= 1; + Some(completion) + } + + /// 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<(), SyncIoError> { + if self.in_flight > 0 { + let (ack, done) = mpsc::sync_channel(1); + self.ops + .send(Op::Barrier(ack)) + .map_err(|_| SyncIoError::WorkerGone)?; + done.recv().map_err(|_| SyncIoError::WorkerGone)?; + } + + if discard { + while self.pop().is_some() {} + } + + Ok(()) + } + + pub fn drain_and_flush(&mut self, discard: bool) -> Result<(), SyncIoError> { + self.drain(discard)?; + + // Sync data out to physical media on host. The worker holds no data of its own, so the + // file descriptor here reaches everything it wrote. + self.file.sync_all().map_err(SyncIoError::SyncAll) + } +} + +impl Drop for SyncFileEngine { + fn drop(&mut self) { + // The worker only exits once it gets here, so everything submitted before is finished. + let _ = self.ops.send(Op::Exit); + if let Some(worker) = self.worker.take() + && worker.join().is_err() + { + error!("The block IO worker thread panicked"); + } + } } diff --git a/src/vmm/src/devices/virtio/block/virtio/request.rs b/src/vmm/src/devices/virtio/block/virtio/request.rs index 8fc83cf43da..2f82639e7d0 100644 --- a/src/vmm/src/devices/virtio/block/virtio/request.rs +++ b/src/vmm/src/devices/virtio/block/virtio/request.rs @@ -397,9 +397,6 @@ impl Request { match res { Ok(block_io::FileEngineOk::Submitted) => ProcessingResult::Submitted, - Ok(block_io::FileEngineOk::Executed(res)) => { - ProcessingResult::Executed(res.req.finish(mem, Ok(res.count), block_metrics)) - } Err(err) => { if err.error.is_throttling_err() { ProcessingResult::Throttled diff --git a/src/vmm/src/devices/virtio/block/virtio/test_utils.rs b/src/vmm/src/devices/virtio/block/virtio/test_utils.rs index e4f23c6a038..b08b620e758 100644 --- a/src/vmm/src/devices/virtio/block/virtio/test_utils.rs +++ b/src/vmm/src/devices/virtio/block/virtio/test_utils.rs @@ -95,14 +95,14 @@ pub fn simulate_queue_event(b: &mut VirtioBlock, maybe_expected_irq: Option { - simulate_queue_event(b, None); - simulate_async_completion_event(b, expected_irq); - } - FileEngine::Sync(_) => { - simulate_queue_event(b, Some(expected_irq)); - } - } + simulate_queue_event(b, None); + simulate_async_completion_event(b, expected_irq); } /// Structure encapsulating the virtq descriptors of a single request to the block device diff --git a/src/vmm/src/seccomp.rs b/src/vmm/src/seccomp.rs index 56e30908a62..e26ad8c20aa 100644 --- a/src/vmm/src/seccomp.rs +++ b/src/vmm/src/seccomp.rs @@ -3,7 +3,7 @@ use std::collections::HashMap; use std::io::Read; -use std::sync::Arc; +use std::sync::{Arc, OnceLock}; use bincode::config; use bincode::config::{Configuration, Fixint, Limit, LittleEndian}; @@ -43,9 +43,27 @@ pub fn get_empty_filters() -> BpfThreadMap { map.insert("vmm".to_string(), Arc::new(vec![])); map.insert("api".to_string(), Arc::new(vec![])); map.insert("vcpu".to_string(), Arc::new(vec![])); + map.insert("block_io".to_string(), Arc::new(vec![])); map } +/// Filter applied by block device IO worker threads, see [`set_block_io_filter`]. +static BLOCK_IO_FILTER: OnceLock> = OnceLock::new(); + +/// Set the filter each block device IO worker thread applies to itself when it starts. +/// +/// Those threads are started whenever a drive is created, which is before the VMM thread's own +/// filter is installed, so they cannot rely on inheriting it. Only the first call has an effect. +pub fn set_block_io_filter(filter: Arc) { + // Ignoring the error is what makes later calls no-ops. + let _ = BLOCK_IO_FILTER.set(filter); +} + +/// The filter set with [`set_block_io_filter`], if any. +pub fn block_io_filter() -> Option> { + BLOCK_IO_FILTER.get().map(|filter| filter.as_slice()) +} + /// Deserialize binary with bpf filters pub fn deserialize_binary(mut reader: R) -> Result { let result: HashMap = bincode::decode_from_std_read(&mut reader, BINCODE_CONFIG)?; From 3920fec768c6dd27cb969fa1f9c7bc173066f6df Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Wed, 23 Sep 2026 21:19:39 -0300 Subject: [PATCH 2/2] vmm: coalesce Sync block completion signals The worker thread signalled the completion eventfd once per request. With the event loop no longer busy doing the I/O itself, it woke up for almost every signal, reaped one or two completions and raised a guest interrupt for each batch, where the inline path used to complete a whole pass over the avail ring with one interrupt. On a fast backing file under load this multiplied the block interrupts the guest had to handle, taking vCPU time from everything else it was doing. Hold the signal back while more requests are queued behind the one that just finished, but only while the previous signal is recent (200us), so a slow backing file still gets one after every request. The worker never blocks waiting for work with a completion left unsignalled, including when the next op is a barrier or a file update. With a guest writing 60k IOPS to a tmpfs-backed drive, block interrupts drop from ~375k to ~100k over 15 s (the inline engine raised ~50k), and completion latency is unaffected on a slow backing file. --- .../src/devices/virtio/block/virtio/io/mod.rs | 30 ++++++++ .../devices/virtio/block/virtio/io/sync_io.rs | 72 ++++++++++++++++++- 2 files changed, 99 insertions(+), 3 deletions(-) diff --git a/src/vmm/src/devices/virtio/block/virtio/io/mod.rs b/src/vmm/src/devices/virtio/block/virtio/io/mod.rs index 5367217daae..184c58376ce 100644 --- a/src/vmm/src/devices/virtio/block/virtio/io/mod.rs +++ b/src/vmm/src/devices/virtio/block/virtio/io/mod.rs @@ -389,6 +389,36 @@ pub mod tests { assert!(engine.pop(&mem).unwrap().is_none()); } + #[test] + fn test_sync_signals_every_completion() { + let mem = create_mem(); + let file = TempFile::new().unwrap().into_file(); + let mut engine = FileEngine::from_file(file, FileEngineType::Sync).unwrap(); + + // Queue ops that complete nothing behind a request: the worker must still signal the + // request's completion before it goes idle. + assert_queued!(engine.write( + 0, + &mem, + GuestAddress(0), + FILE_LEN, + PendingRequest::default() + )); + engine + .update_file_path(TempFile::new().unwrap().into_file()) + .unwrap(); + + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while engine.completion_evt().read().is_err() { + assert!( + std::time::Instant::now() < deadline, + "completion never signalled" + ); + std::thread::sleep(std::time::Duration::from_millis(1)); + } + assert_eq!(engine.pop(&mem).unwrap().unwrap().result.unwrap(), FILE_LEN); + } + #[test] fn test_sync_throttling() { let mem = create_mem(); diff --git a/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs b/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs index 5495930a586..3254e5254bb 100644 --- a/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs +++ b/src/vmm/src/devices/virtio/block/virtio/io/sync_io.rs @@ -5,13 +5,14 @@ //! //! The I/O itself is done with blocking system calls, but not on the thread that submits it: each //! engine owns a worker thread that performs the requests one at a time, in submission order, and -//! reports each completion through an eventfd, the same way the async engine does. This keeps a +//! reports completions through an eventfd, the same way the async engine does. This keeps a //! slow backing file from stalling the event loop, and with it every other device serviced there. use std::fs::File; use std::io::{Seek, SeekFrom, Write}; use std::sync::mpsc; use std::thread; +use std::time::{Duration, Instant}; use vm_memory::{GuestMemoryError, ReadVolatile, WriteVolatile}; use vmm_sys_util::eventfd::EventFd; @@ -168,7 +169,21 @@ fn run_worker( panic!("Failed to set the requested seccomp filters on the block IO worker: {err}"); } - for op in ops { + let mut signal = CompletionSignal::new(completion_evt); + let mut next = None; + loop { + let op = match next.take() { + Some(op) => op, + None => { + // Never go to sleep on a completion the device has not been told about. + signal.flush(); + match ops.recv() { + Ok(op) => op, + Err(mpsc::RecvError) => break, + } + } + }; + let completion = match op { Op::Io { io, req } => SyncCompletion { req, @@ -189,9 +204,60 @@ fn run_worker( if completions.send(completion).is_err() { break; } - if let Err(err) = completion_evt.write(1) { + signal.pending = true; + + // Hold the signal back while more requests are queued, so that the device handles a + // burst of completions at once instead of raising an interrupt for each. Only while the + // previous signal is recent, though: a slow backing file gets one after every request. + match ops.try_recv() { + Ok(op) => { + next = Some(op); + signal.maybe_flush(); + } + Err(_) => signal.flush(), + } + } + signal.flush(); +} + +/// How long a finished request may wait for the requests queued behind it before the device is +/// told about it. +const COMPLETION_SIGNAL_DELAY: Duration = Duration::from_micros(200); + +/// Coalesces the completion eventfd writes of the worker thread. +#[derive(Debug)] +struct CompletionSignal { + evt: EventFd, + pending: bool, + last: Instant, +} + +impl CompletionSignal { + fn new(evt: EventFd) -> Self { + CompletionSignal { + evt, + pending: false, + last: Instant::now(), + } + } + + /// Signal pending completions if the last signal is old enough. + fn maybe_flush(&mut self) { + if self.last.elapsed() >= COMPLETION_SIGNAL_DELAY { + self.flush(); + } + } + + /// Signal pending completions. + fn flush(&mut self) { + if !self.pending { + return; + } + if let Err(err) = self.evt.write(1) { error!("Failed to signal block IO completion: {:?}", err); } + self.pending = false; + self.last = Instant::now(); } }