diff --git a/vm-virtio/src/block.rs b/vm-virtio/src/block.rs index 50b808223..de6169294 100755 --- a/vm-virtio/src/block.rs +++ b/vm-virtio/src/block.rs @@ -22,25 +22,22 @@ use std::cmp; use std::convert::TryInto; use std::fs::{File, Metadata}; use std::io::{self, Read, Seek, SeekFrom, Write}; +use std::ops::DerefMut; use std::os::linux::fs::MetadataExt; use std::os::unix::io::{AsRawFd, RawFd}; use std::path::PathBuf; use std::result; use std::slice; use std::sync::atomic::{AtomicBool, Ordering}; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use std::thread; use virtio_bindings::bindings::virtio_blk::*; use vm_device::{Migratable, MigratableError, Pausable, Snapshotable}; -use vm_memory::{Bytes, GuestAddress, GuestMemory, GuestMemoryError, GuestMemoryMmap}; +use vm_memory::{ByteValued, Bytes, GuestAddress, GuestMemory, GuestMemoryError, GuestMemoryMmap}; use vmm_sys_util::{eventfd::EventFd, seek_hole::SeekHole, write_zeroes::PunchHole}; -const CONFIG_SPACE_SIZE: usize = 8; const SECTOR_SHIFT: u8 = 9; pub const SECTOR_SIZE: u64 = (0x01 as u64) << SECTOR_SHIFT; -const QUEUE_SIZE: u16 = 256; -const NUM_QUEUES: usize = 1; -const QUEUE_SIZES: &[u16] = &[QUEUE_SIZE]; // New descriptors are pending on the virtio queue. const QUEUE_AVAIL_EVENT: DeviceEventT = 0; @@ -592,9 +589,9 @@ impl Request { } struct BlockEpollHandler { - queues: Vec, + queue: Queue, mem: Arc>, - disk_image: T, + disk_image: Arc>, disk_nsectors: u64, interrupt_cb: Arc, disk_image_id: Vec, @@ -603,18 +600,20 @@ struct BlockEpollHandler { } impl BlockEpollHandler { - fn process_queue(&mut self, queue_index: usize) -> bool { - let queue = &mut self.queues[queue_index]; + fn process_queue(&mut self) -> bool { + let queue = &mut self.queue; - let mut used_desc_heads = [(0, 0); QUEUE_SIZE as usize]; + let mut used_desc_heads = Vec::new(); let mut used_count = 0; let mem = self.mem.load(); for avail_desc in queue.iter(&mem) { let len; match Request::parse(&avail_desc, &mem) { Ok(request) => { + let mut disk_image_locked = self.disk_image.lock().unwrap(); + let mut disk_image = disk_image_locked.deref_mut(); let status = match request.execute( - &mut self.disk_image, + &mut disk_image, self.disk_nsectors, &mem, &self.disk_image_id, @@ -638,19 +637,19 @@ impl BlockEpollHandler { len = 0; } } - used_desc_heads[used_count] = (avail_desc.index, len); + used_desc_heads.push((avail_desc.index, len)); used_count += 1; } - for &(desc_index, len) in &used_desc_heads[..used_count] { + for &(desc_index, len) in used_desc_heads.iter() { queue.add_used(&mem, desc_index, len); } used_count > 0 } - fn signal_used_queue(&self, queue_index: usize) -> result::Result<(), DeviceError> { + fn signal_used_queue(&self) -> result::Result<(), DeviceError> { self.interrupt_cb - .trigger(&VirtioInterruptType::Queue, Some(&self.queues[queue_index])) + .trigger(&VirtioInterruptType::Queue, Some(&self.queue)) .map_err(|e| { error!("Failed to signal used queue: {:?}", e); DeviceError::FailedSignalingUsedQueue(e) @@ -660,16 +659,15 @@ impl BlockEpollHandler { #[allow(dead_code)] fn update_disk_image( &mut self, - disk_image: T, + mut disk_image: T, disk_path: &PathBuf, ) -> result::Result<(), DeviceError> { - self.disk_image = disk_image; - self.disk_nsectors = self - .disk_image + self.disk_nsectors = disk_image .seek(SeekFrom::End(0)) .map_err(DeviceError::IoError)? / SECTOR_SIZE; self.disk_image_id = build_disk_image_id(disk_path); + self.disk_image = Arc::new(Mutex::new(disk_image)); Ok(()) } @@ -733,8 +731,8 @@ impl BlockEpollHandler { if let Err(e) = queue_evt.read() { error!("Failed to get queue event: {:?}", e); break 'epoll; - } else if self.process_queue(0) { - if let Err(e) = self.signal_used_queue(0) { + } else if self.process_queue() { + if let Err(e) = self.signal_used_queue() { error!("Failed to signal used queue: {:?}", e); break 'epoll; } @@ -764,32 +762,57 @@ impl BlockEpollHandler { } } +#[derive(Copy, Clone, Debug, Default)] +#[repr(C, packed)] +pub struct VirtioBlockGeometry { + pub cylinders: u16, + pub heads: u8, + pub sectors: u8, +} + +unsafe impl ByteValued for VirtioBlockGeometry {} + +#[derive(Copy, Clone, Debug, Default)] +#[repr(C, packed)] +pub struct VirtioBlockConfig { + pub capacity: u64, + pub size_max: u32, + pub seg_max: u32, + pub geometry: VirtioBlockGeometry, + pub blk_size: u32, + pub physical_block_exp: u8, + pub alignment_offset: u8, + pub min_io_size: u16, + pub opt_io_size: u32, + pub wce: u8, + unused: u8, + pub num_queues: u16, + pub max_discard_sectors: u32, + pub max_discard_seg: u32, + pub discard_sector_alignment: u32, + pub max_write_zeroes_sectors: u32, + pub max_write_zeroes_seg: u32, + pub write_zeroes_may_unmap: u8, + unused1: [u8; 3], +} + +unsafe impl ByteValued for VirtioBlockConfig {} + /// Virtio device for exposing block level read/write operations on a host file. pub struct Block { kill_evt: Option, - disk_image: Option, + disk_image: Arc>, disk_path: PathBuf, disk_nsectors: u64, avail_features: u64, acked_features: u64, - config_space: Vec, - queue_evt: Option, + config: VirtioBlockConfig, + queue_evts: Option>, interrupt_cb: Option>, epoll_threads: Option>>>, pause_evt: Option, paused: Arc, -} - -pub fn build_config_space(disk_size: u64) -> Vec { - // We only support disk size, which uses the first two words of the configuration space. - // If the image is not a multiple of the sector size, the tail bits are not exposed. - // The config space is little endian. - let mut config = Vec::with_capacity(CONFIG_SPACE_SIZE); - let num_sectors = disk_size >> SECTOR_SHIFT; - for i in 0..8 { - config.push((num_sectors >> (8 * i)) as u8); - } - config + queue_size: Vec, } impl Block { @@ -801,6 +824,8 @@ impl Block { disk_path: PathBuf, is_disk_read_only: bool, iommu: bool, + num_queues: usize, + queue_size: u16, ) -> io::Result> { let disk_size = disk_image.seek(SeekFrom::End(0))? as u64; if disk_size % SECTOR_SIZE != 0 { @@ -819,21 +844,33 @@ impl Block { if is_disk_read_only { avail_features |= 1u64 << VIRTIO_BLK_F_RO; + } + + let disk_nsectors = disk_size / SECTOR_SIZE; + let mut config = VirtioBlockConfig { + capacity: disk_nsectors, + ..Default::default() }; + if num_queues > 1 { + avail_features |= 1u64 << VIRTIO_BLK_F_MQ; + config.num_queues = num_queues as u16; + } + Ok(Block { kill_evt: None, - disk_image: Some(disk_image), + disk_image: Arc::new(Mutex::new(disk_image)), disk_path, - disk_nsectors: disk_size / SECTOR_SIZE, + disk_nsectors, avail_features, acked_features: 0u64, - config_space: build_config_space(disk_size), - queue_evt: None, + config, + queue_evts: None, interrupt_cb: None, epoll_threads: None, pause_evt: None, paused: Arc::new(AtomicBool::new(false)), + queue_size: vec![queue_size; num_queues], }) } } @@ -853,7 +890,7 @@ impl VirtioDevice for Block { } fn queue_max_sizes(&self) -> &[u16] { - QUEUE_SIZES + self.queue_size.as_slice() } fn features(&self, page: u32) -> u32 { @@ -891,26 +928,28 @@ impl VirtioDevice for Block { } fn read_config(&self, offset: u64, mut data: &mut [u8]) { - let config_len = self.config_space.len() as u64; + let config_slice = self.config.as_slice(); + let config_len = config_slice.len() as u64; if offset >= config_len { error!("Failed to read config space"); return; } if let Some(end) = offset.checked_add(data.len() as u64) { // This write can't fail, offset and end are checked against config_len. - data.write_all(&self.config_space[offset as usize..cmp::min(end, config_len) as usize]) + data.write_all(&config_slice[offset as usize..cmp::min(end, config_len) as usize]) .unwrap(); } } fn write_config(&mut self, offset: u64, data: &[u8]) { + let config_slice = self.config.as_mut_slice(); let data_len = data.len() as u64; - let config_len = self.config_space.len() as u64; + let config_len = config_slice.len() as u64; if offset + data_len > config_len { error!("Failed to write config space"); return; } - let (_, right) = self.config_space.split_at_mut(offset as usize); + let (_, right) = config_slice.split_at_mut(offset as usize); right.copy_from_slice(&data[..]); } @@ -918,13 +957,13 @@ impl VirtioDevice for Block { &mut self, mem: Arc>, interrupt_cb: Arc, - queues: Vec, + mut queues: Vec, mut queue_evts: Vec, ) -> ActivateResult { - if queues.len() != NUM_QUEUES || queue_evts.len() != NUM_QUEUES { + if queues.len() != self.queue_size.len() || queue_evts.len() != self.queue_size.len() { error!( "Cannot perform activate. Expected {} queue(s), got {}", - NUM_QUEUES, + self.queue_size.len(), queues.len() ); return Err(ActivateError::BadActivate); @@ -947,35 +986,45 @@ impl VirtioDevice for Block { })?; self.pause_evt = Some(self_pause_evt); - if let Some(disk_image) = self.disk_image.clone() { - let disk_image_id = build_disk_image_id(&self.disk_path); - - // Save the interrupt EventFD as we need to return it on reset - // but clone it to pass into the thread. - self.interrupt_cb = Some(interrupt_cb); - let interrupt_cb = self.interrupt_cb.as_ref().unwrap().clone(); + let disk_image_id = build_disk_image_id(&self.disk_path); + let mut tmp_queue_evts: Vec = Vec::new(); + for queue_evt in queue_evts.iter() { // Save the queue EventFD as we need to return it on reset // but clone it to pass into the thread. - self.queue_evt = Some(queue_evts.remove(0)); - let queue_evt = self.queue_evt.as_ref().unwrap().try_clone().map_err(|e| { + tmp_queue_evts.push(queue_evt.try_clone().map_err(|e| { error!("failed to clone queue EventFd: {}", e); ActivateError::BadActivate - })?; + })?); + } + self.queue_evts = Some(tmp_queue_evts); + let mut tmp_queue_evts: Vec = Vec::new(); + for queue_evt in queue_evts.iter() { + // Save the queue EventFD as we need to return it on reset + // but clone it to pass into the thread. + tmp_queue_evts.push(queue_evt.try_clone().map_err(|e| { + error!("failed to clone queue EventFd: {}", e); + ActivateError::BadActivate + })?); + } + self.queue_evts = Some(tmp_queue_evts); + + let mut epoll_threads = Vec::new(); + for _ in 0..self.queue_size.len() { let mut handler = BlockEpollHandler { - queues, - mem, - disk_image, + queue: queues.remove(0), + mem: mem.clone(), + disk_image: self.disk_image.clone(), disk_nsectors: self.disk_nsectors, - interrupt_cb, - disk_image_id, - kill_evt, - pause_evt, + interrupt_cb: interrupt_cb.clone(), + disk_image_id: disk_image_id.clone(), + kill_evt: kill_evt.try_clone().unwrap(), + pause_evt: pause_evt.try_clone().unwrap(), }; + let queue_evt = queue_evts.remove(0); let paused = self.paused.clone(); - let mut epoll_threads = Vec::new(); thread::Builder::new() .name("virtio_blk".to_string()) .spawn(move || handler.run(queue_evt, paused)) @@ -984,12 +1033,15 @@ impl VirtioDevice for Block { error!("failed to clone the virtio-blk epoll thread: {}", e); ActivateError::BadActivate })?; - - self.epoll_threads = Some(epoll_threads); - - return Ok(()); } - Err(ActivateError::BadActivate) + + // Save the interrupt EventFD as we need to return it on reset + // but clone it to pass into the thread. + self.interrupt_cb = Some(interrupt_cb); + + self.epoll_threads = Some(epoll_threads); + + Ok(()) } fn reset(&mut self) -> Option<(Arc, Vec)> { @@ -1006,7 +1058,7 @@ impl VirtioDevice for Block { // Return the interrupt and queue EventFDs Some(( self.interrupt_cb.take().unwrap(), - vec![self.queue_evt.take().unwrap()], + self.queue_evts.take().unwrap(), )) } } diff --git a/vmm/src/device_manager.rs b/vmm/src/device_manager.rs index e8a30ef37..e45ad5815 100644 --- a/vmm/src/device_manager.rs +++ b/vmm/src/device_manager.rs @@ -991,6 +991,8 @@ impl DeviceManager { disk_cfg.path.clone(), disk_cfg.readonly, disk_cfg.iommu, + disk_cfg.num_queues, + disk_cfg.queue_size, ) .map_err(DeviceManagerError::CreateVirtioBlock)?; @@ -1010,6 +1012,8 @@ impl DeviceManager { disk_cfg.path.clone(), disk_cfg.readonly, disk_cfg.iommu, + disk_cfg.num_queues, + disk_cfg.queue_size, ) .map_err(DeviceManagerError::CreateVirtioBlock)?;