1use std::cell::RefCell;
6use std::collections::BTreeMap;
7use std::collections::BTreeSet;
8use std::io;
9use std::io::Write;
10use std::mem::size_of;
11use std::rc::Rc;
12use std::result;
13use std::sync::atomic::AtomicU64;
14use std::sync::atomic::Ordering;
15use std::sync::Arc;
16use std::time::Duration;
17
18use anyhow::Context;
19use base::debug;
20use base::error;
21use base::info;
22use base::warn;
23use base::AsRawDescriptor;
24use base::Error as SysError;
25use base::Event;
26use base::RawDescriptor;
27use base::Result as SysResult;
28use base::Timer;
29use base::Tube;
30use base::TubeError;
31use base::WorkerThread;
32use cros_async::sync::RwLock as AsyncRwLock;
33use cros_async::AsyncError;
34use cros_async::AsyncTube;
35use cros_async::EventAsync;
36use cros_async::Executor;
37use cros_async::ExecutorKind;
38use cros_async::IoOptions;
39use cros_async::TimerAsync;
40use data_model::Le16;
41use data_model::Le32;
42use data_model::Le64;
43use devices::virtio::async_utils;
44use devices::virtio::copy_config;
45use devices::virtio::device_constants::block::virtio_blk_config;
46use devices::virtio::device_constants::block::virtio_blk_discard_write_zeroes;
47use devices::virtio::device_constants::block::virtio_blk_req_header;
48use devices::virtio::device_constants::block::VIRTIO_BLK_DISCARD_WRITE_ZEROES_FLAG_UNMAP;
49use devices::virtio::device_constants::block::VIRTIO_BLK_F_BLK_SIZE;
50use devices::virtio::device_constants::block::VIRTIO_BLK_F_DISCARD;
51use devices::virtio::device_constants::block::VIRTIO_BLK_F_FLUSH;
52use devices::virtio::device_constants::block::VIRTIO_BLK_F_MQ;
53use devices::virtio::device_constants::block::VIRTIO_BLK_F_RO;
54use devices::virtio::device_constants::block::VIRTIO_BLK_F_SEG_MAX;
55use devices::virtio::device_constants::block::VIRTIO_BLK_F_WRITE_ZEROES;
56use devices::virtio::device_constants::block::VIRTIO_BLK_S_IOERR;
57use devices::virtio::device_constants::block::VIRTIO_BLK_S_OK;
58use devices::virtio::device_constants::block::VIRTIO_BLK_S_UNSUPP;
59use devices::virtio::device_constants::block::VIRTIO_BLK_T_DISCARD;
60use devices::virtio::device_constants::block::VIRTIO_BLK_T_FLUSH;
61use devices::virtio::device_constants::block::VIRTIO_BLK_T_GET_ID;
62use devices::virtio::device_constants::block::VIRTIO_BLK_T_IN;
63use devices::virtio::device_constants::block::VIRTIO_BLK_T_OUT;
64use devices::virtio::device_constants::block::VIRTIO_BLK_T_WRITE_ZEROES;
65use devices::virtio::DescriptorChain;
66use devices::virtio::DeviceType;
67use devices::virtio::Interrupt;
68use devices::virtio::Queue;
69use devices::virtio::Reader;
70use devices::virtio::VirtioDevice;
71use devices::virtio::Writer;
72use devices::PciAddress;
73use disk::AsyncDisk;
74use disk::DiskFile;
75use futures::channel::mpsc;
76use futures::channel::oneshot;
77use futures::pin_mut;
78use futures::stream::FuturesUnordered;
79use futures::stream::StreamExt;
80use futures::FutureExt;
81use remain::sorted;
82use snapshot::AnySnapshot;
83use thiserror::Error as ThisError;
84use virtio_sys::virtio_config::VIRTIO_F_RING_PACKED;
85use vm_control::DiskControlCommand;
86use vm_control::DiskControlResult;
87use vm_memory::GuestMemory;
88use zerocopy::IntoBytes;
89
90use crate::sys::*;
91use crate::DiskOption;
92
93const DEFAULT_QUEUE_SIZE: u16 = 256;
94const DEFAULT_NUM_QUEUES: u16 = 16;
95
96const SECTOR_SHIFT: u8 = 9;
97const SECTOR_SIZE: u64 = 0x01 << SECTOR_SHIFT;
98
99const MAX_DISCARD_SECTORS: u32 = u32::MAX;
100const MAX_WRITE_ZEROES_SECTORS: u32 = u32::MAX;
101const MAX_DISCARD_SEG: u32 = 32;
103const MAX_WRITE_ZEROES_SEG: u32 = 32;
104const DISCARD_SECTOR_ALIGNMENT: u32 = 128;
107
108#[sorted]
109#[derive(ThisError, Debug)]
110enum ExecuteError {
111 #[error("failed to copy ID string: {0}")]
112 CopyId(io::Error),
113 #[error("failed to perform discard or write zeroes; sector={sector} num_sectors={num_sectors} flags={flags}; {ioerr:?}")]
114 DiscardWriteZeroes {
115 ioerr: Option<disk::Error>,
116 sector: u64,
117 num_sectors: u32,
118 flags: u32,
119 },
120 #[error("failed to flush: {0}")]
121 Flush(disk::Error),
122 #[error("not enough space in descriptor chain to write status")]
123 MissingStatus,
124 #[error("out of range")]
125 OutOfRange,
126 #[error("failed to read message: {0}")]
127 Read(io::Error),
128 #[error("io error reading {length} bytes from sector {sector}: {desc_error}")]
129 ReadIo {
130 length: usize,
131 sector: u64,
132 desc_error: disk::Error,
133 },
134 #[error("read only; request_type={request_type}")]
135 ReadOnly { request_type: u32 },
136 #[error("failed to recieve command message: {0}")]
137 ReceivingCommand(TubeError),
138 #[error("failed to send command response: {0}")]
139 SendingResponse(TubeError),
140 #[error("couldn't reset the timer: {0}")]
141 TimerReset(base::Error),
142 #[error("too many segments: {0} > {0}")]
143 TooManySegments(usize, usize),
144 #[error("unsupported ({0})")]
145 Unsupported(u32),
146 #[error("io error writing {length} bytes from sector {sector}: {desc_error}")]
147 WriteIo {
148 length: usize,
149 sector: u64,
150 desc_error: disk::Error,
151 },
152 #[error("failed to write request status: {0}")]
153 WriteStatus(io::Error),
154}
155
156enum LogLevel {
157 Debug,
158 Error,
159}
160
161impl ExecuteError {
162 fn status(&self) -> u8 {
163 match self {
164 ExecuteError::CopyId(_) => VIRTIO_BLK_S_IOERR,
165 ExecuteError::DiscardWriteZeroes { .. } => VIRTIO_BLK_S_IOERR,
166 ExecuteError::Flush(_) => VIRTIO_BLK_S_IOERR,
167 ExecuteError::MissingStatus => VIRTIO_BLK_S_IOERR,
168 ExecuteError::OutOfRange => VIRTIO_BLK_S_IOERR,
169 ExecuteError::Read(_) => VIRTIO_BLK_S_IOERR,
170 ExecuteError::ReadIo { .. } => VIRTIO_BLK_S_IOERR,
171 ExecuteError::ReadOnly { .. } => VIRTIO_BLK_S_IOERR,
172 ExecuteError::ReceivingCommand(_) => VIRTIO_BLK_S_IOERR,
173 ExecuteError::SendingResponse(_) => VIRTIO_BLK_S_IOERR,
174 ExecuteError::TimerReset(_) => VIRTIO_BLK_S_IOERR,
175 ExecuteError::TooManySegments(_, _) => VIRTIO_BLK_S_IOERR,
176 ExecuteError::WriteIo { .. } => VIRTIO_BLK_S_IOERR,
177 ExecuteError::WriteStatus(_) => VIRTIO_BLK_S_IOERR,
178 ExecuteError::Unsupported(_) => VIRTIO_BLK_S_UNSUPP,
179 }
180 }
181
182 fn log_level(&self) -> LogLevel {
183 match self {
184 ExecuteError::ReadIo { .. }
186 | ExecuteError::WriteIo { .. }
187 | ExecuteError::Flush { .. }
188 | ExecuteError::DiscardWriteZeroes { .. } => LogLevel::Debug,
189 _ => LogLevel::Error,
191 }
192 }
193}
194
195#[sorted]
198#[derive(ThisError, Debug)]
199enum ControlError {
200 #[error("failed to fdatasync the disk: {0}")]
201 FdatasyncDisk(disk::Error),
202 #[error("couldn't get a value from a timer for flushing: {0}")]
203 FlushTimer(AsyncError),
204}
205
206const ID_LEN: usize = 20;
208
209type BlockId = [u8; ID_LEN];
213
214struct DiskState {
216 disk_image: Box<dyn AsyncDisk>,
217 read_only: bool,
218 sparse: bool,
219 id: BlockId,
220 dontcache_read: bool,
221 dontcache_write: bool,
222 worker_shared_state: Arc<AsyncRwLock<WorkerSharedState>>,
225}
226
227struct WorkerSharedState {
229 disk_size: Arc<AtomicU64>,
230}
231
232async fn process_one_request(
233 avail_desc: &mut DescriptorChain,
234 disk_state: &AsyncRwLock<DiskState>,
235 flush_timer: &RefCell<TimerAsync<Timer>>,
236 flush_timer_armed: &RefCell<bool>,
237) -> result::Result<usize, ExecuteError> {
238 let reader = &mut avail_desc.reader;
239 let writer = &mut avail_desc.writer;
240
241 let available_bytes = writer.available_bytes();
245 let status_offset = available_bytes
246 .checked_sub(1)
247 .ok_or(ExecuteError::MissingStatus)?;
248 let mut status_writer = writer.split_at(status_offset);
249
250 let status = match BlockAsync::execute_request(
251 reader,
252 writer,
253 disk_state,
254 flush_timer,
255 flush_timer_armed,
256 )
257 .await
258 {
259 Ok(()) => VIRTIO_BLK_S_OK,
260 Err(e) => {
261 match e.log_level() {
262 LogLevel::Debug => debug!("failed executing disk request: {:#}", e),
263 LogLevel::Error => error!("failed executing disk request: {:#}", e),
264 }
265 e.status()
266 }
267 };
268
269 status_writer
270 .write_all(&[status])
271 .map_err(ExecuteError::WriteStatus)?;
272 Ok(available_bytes)
273}
274
275async fn process_one_chain(
277 queue: &RefCell<Queue>,
278 mut avail_desc: DescriptorChain,
279 disk_state: &AsyncRwLock<DiskState>,
280 flush_timer: &RefCell<TimerAsync<Timer>>,
281 flush_timer_armed: &RefCell<bool>,
282) {
283 let len = match process_one_request(&mut avail_desc, disk_state, flush_timer, flush_timer_armed)
284 .await
285 {
286 Ok(len) => len,
287 Err(e) => {
288 error!("block: failed to handle request: {:#}", e);
289 0
290 }
291 };
292
293 let mut queue = queue.borrow_mut();
294 queue.add_used_with_bytes_written(avail_desc, len as u32);
295 queue.trigger_interrupt();
296}
297
298async fn handle_queue(
302 disk_state: Rc<AsyncRwLock<DiskState>>,
303 queue: Queue,
304 evt: EventAsync,
305 flush_timer: Rc<RefCell<TimerAsync<Timer>>>,
306 flush_timer_armed: Rc<RefCell<bool>>,
307 mut stop_rx: oneshot::Receiver<()>,
308) -> Queue {
309 let queue = RefCell::new(queue);
310 let mut background_tasks = FuturesUnordered::new();
311 let evt_future = futures::future::Either::Left(std::future::ready(Ok(0))).fuse();
312 pin_mut!(evt_future);
313 loop {
314 futures::select! {
321 _ = background_tasks.next() => continue,
322 res = evt_future => {
323 evt_future.set(futures::future::Either::Right(evt.next_val()).fuse());
324 if let Err(e) = res {
325 error!("Failed to read the next queue event: {:#}", e);
326 continue;
327 }
328 }
329 _ = stop_rx => {
330 background_tasks.collect::<()>().await;
333 return queue.into_inner();
334 }
335 };
336 while let Some(descriptor_chain) = queue.borrow_mut().pop() {
337 background_tasks.push(process_one_chain(
338 &queue,
339 descriptor_chain,
340 &disk_state,
341 &flush_timer,
342 &flush_timer_armed,
343 ));
344 }
345 }
346}
347
348async fn handle_command_tube(
349 command_tube: &Option<AsyncTube>,
350 interrupt: &RefCell<Option<Interrupt>>,
351 disk_state: Rc<AsyncRwLock<DiskState>>,
352) -> Result<(), ExecuteError> {
353 let command_tube = match command_tube {
354 Some(c) => c,
355 None => {
356 futures::future::pending::<()>().await;
357 return Ok(());
358 }
359 };
360 loop {
361 match command_tube.next().await {
362 Ok(command) => {
363 let resp = match command {
364 DiskControlCommand::Resize { new_size } => resize(&disk_state, new_size).await,
365 };
366
367 let resp_clone = resp.clone();
368 command_tube
369 .send(resp_clone)
370 .await
371 .map_err(ExecuteError::SendingResponse)?;
372 if let DiskControlResult::Ok = resp {
373 if let Some(interrupt) = &*interrupt.borrow() {
374 interrupt.signal_config_changed();
375 }
376 }
377 }
378 Err(e) => return Err(ExecuteError::ReceivingCommand(e)),
379 }
380 }
381}
382
383async fn resize(disk_state: &AsyncRwLock<DiskState>, new_size: u64) -> DiskControlResult {
384 let disk_state = disk_state.lock().await;
387 let worker_shared_state = Arc::clone(&disk_state.worker_shared_state);
389 let worker_shared_state = worker_shared_state.lock().await;
390
391 if disk_state.read_only {
392 error!("Attempted to resize read-only block device");
393 return DiskControlResult::Err(SysError::new(libc::EROFS));
394 }
395
396 info!("Resizing block device to {} bytes", new_size);
397
398 if let Err(e) = disk_state.disk_image.set_len(new_size) {
399 error!("Resizing disk failed! {:#}", e);
400 return DiskControlResult::Err(SysError::new(libc::EIO));
401 }
402
403 if !disk_state.sparse {
405 if let Err(e) = disk_state.disk_image.allocate(0, new_size) {
406 error!("Allocating disk space after resize failed! {:#}", e);
407 return DiskControlResult::Err(SysError::new(libc::EIO));
408 }
409 }
410
411 if let Ok(new_disk_size) = disk_state.disk_image.get_len() {
412 worker_shared_state
413 .disk_size
414 .store(new_disk_size, Ordering::Release);
415 }
416 DiskControlResult::Ok
417}
418
419async fn flush_disk(
421 disk_state: Rc<AsyncRwLock<DiskState>>,
422 timer: TimerAsync<Timer>,
423 armed: Rc<RefCell<bool>>,
424) -> Result<(), ControlError> {
425 loop {
426 timer.wait().await.map_err(ControlError::FlushTimer)?;
427 if !*armed.borrow() {
428 continue;
429 }
430
431 *armed.borrow_mut() = false;
434
435 disk_state
436 .read_lock()
437 .await
438 .disk_image
439 .fdatasync()
440 .await
441 .map_err(ControlError::FdatasyncDisk)?;
442 }
443}
444
445enum WorkerCmd {
446 StartQueue {
447 index: usize,
448 queue: Queue,
449 },
450 StopQueue {
451 index: usize,
452 response_tx: oneshot::Sender<Option<Queue>>,
455 },
456 AbortQueues {
459 response_tx: oneshot::Sender<()>,
461 },
462 FlushDisk {
464 response_tx: oneshot::Sender<()>,
466 },
467}
468
469async fn run_worker(
476 ex: &Executor,
477 disk_state: &Rc<AsyncRwLock<DiskState>>,
478 control_tube: &Option<AsyncTube>,
479 mut worker_rx: mpsc::UnboundedReceiver<WorkerCmd>,
480 kill_evt: Event,
481) -> anyhow::Result<()> {
482 let timer = Timer::new().expect("Failed to create a timer");
484 let flush_timer_armed = Rc::new(RefCell::new(false));
485
486 let control_interrupt = RefCell::new(None);
488 let control = handle_command_tube(control_tube, &control_interrupt, disk_state.clone()).fuse();
489 pin_mut!(control);
490
491 let flush_timer = Rc::new(RefCell::new(
493 TimerAsync::new(
494 timer.try_clone().expect("Failed to clone flush_timer"),
496 ex,
497 )
498 .expect("Failed to create an async timer"),
499 ));
500
501 let flush_timer2 = TimerAsync::new(timer, ex).expect("Failed to create an async timer");
503 let disk_flush = flush_disk(disk_state.clone(), flush_timer2, flush_timer_armed.clone()).fuse();
504 pin_mut!(disk_flush);
505
506 let kill = async_utils::await_and_exit(ex, kill_evt).fuse();
508 pin_mut!(kill);
509
510 let mut queue_handlers = FuturesUnordered::new();
512 let mut queue_handler_stop_fns = std::collections::BTreeMap::new();
514
515 loop {
516 futures::select! {
517 _ = queue_handlers.next() => continue,
518 r = disk_flush => return r.context("failed to flush a disk"),
519 r = control => return r.context("failed to handle a control request"),
520 r = kill => return r.context("failed to wait on the kill event"),
521 worker_cmd = worker_rx.next() => {
522 match worker_cmd {
523 None => anyhow::bail!("worker control channel unexpectedly closed"),
524 Some(WorkerCmd::StartQueue{index, queue}) => {
525 if control_interrupt.borrow().is_none() {
526 *control_interrupt.borrow_mut() = Some(queue.interrupt().clone());
527 }
528
529 let (tx, rx) = oneshot::channel();
530 let kick_evt = queue.event().try_clone().expect("Failed to clone queue event");
531 let (handle_queue_future, remote_handle) = handle_queue(
532 Rc::clone(disk_state),
533 queue,
534 EventAsync::new(kick_evt, ex).expect("Failed to create async event for queue"),
535 Rc::clone(&flush_timer),
536 Rc::clone(&flush_timer_armed),
537 rx,
538 ).remote_handle();
539 let old_stop_fn = queue_handler_stop_fns.insert(index, move || {
540 tx.send(()).unwrap_or_else(|_| panic!("queue handler channel closed early"));
542 remote_handle
544 });
545
546 if let Some(stop_fn) = old_stop_fn {
549 warn!("Starting new queue handler without stopping old handler");
550 let mut fut = stop_fn().fuse();
553 loop {
554 futures::select! {
555 _ = queue_handlers.next() => continue,
556 _queue = fut => break,
557 }
558 }
559 }
560
561 queue_handlers.push(handle_queue_future);
562 }
563 Some(WorkerCmd::StopQueue{index, response_tx}) => {
564 match queue_handler_stop_fns.remove(&index) {
565 Some(stop_fn) => {
566 let mut fut = stop_fn().fuse();
574 let queue = loop {
575 futures::select! {
576 _ = queue_handlers.next() => continue,
577 queue = fut => break queue,
578 }
579 };
580
581 if queue_handlers.is_empty() {
585 *control_interrupt.borrow_mut() = None;
586 }
587
588 let _ = response_tx.send(Some(queue));
589 }
590 None => { let _ = response_tx.send(None); },
591 }
592
593 }
594 Some(WorkerCmd::AbortQueues{response_tx}) => {
595 queue_handlers.clear();
596 queue_handler_stop_fns.clear();
597
598 *control_interrupt.borrow_mut() = None;
599
600 let _ = response_tx.send(());
601 }
602 Some(WorkerCmd::FlushDisk{response_tx}) => {
603 let disk = disk_state.read_lock().await;
604 if let Err(e) = disk.disk_image.flush().await {
605 error!("failed to flush disk image: {:#}", e);
606 }
607 let _ = response_tx.send(());
608 }
609 }
610 }
611 };
612 }
613}
614
615pub struct BlockAsync {
617 boot_index: Option<usize>,
620 disk_image: Option<Box<dyn DiskFile>>,
622 disk_size: Arc<AtomicU64>,
623 avail_features: u64,
624 read_only: bool,
625 sparse: bool,
626 seg_max: u32,
627 block_size: u32,
628 id: BlockId,
629 control_tube: Option<Tube>,
630 queue_sizes: Vec<u16>,
631 pub(super) executor_kind: ExecutorKind,
632 worker_threads: BTreeMap<usize, (WorkerThread<()>, mpsc::UnboundedSender<WorkerCmd>)>,
640 shared_state: Arc<AsyncRwLock<WorkerSharedState>>,
641 worker_per_queue: bool,
643 activated_queues: BTreeSet<usize>,
647 #[cfg(windows)]
648 pub(super) io_concurrency: u32,
649 pci_address: Option<PciAddress>,
650 dontcache_read: bool,
651 dontcache_write: bool,
652}
653
654impl BlockAsync {
655 pub fn new(
657 base_features: u64,
658 disk_image: Box<dyn DiskFile>,
659 disk_option: &DiskOption,
660 control_tube: Option<Tube>,
661 queue_size: Option<u16>,
662 num_queues: Option<u16>,
663 ) -> SysResult<BlockAsync> {
664 let read_only = disk_option.read_only;
665 let sparse = disk_option.sparse;
666 let block_size = disk_option.block_size;
667 let packed_queue = disk_option.packed_queue;
668 let id = disk_option.id.unwrap_or_default();
669 let mut worker_per_queue = disk_option.multiple_workers;
670 if worker_per_queue && disk_image.try_clone().is_err() {
672 base::warn!("multiple workers requested, but not supported by disk image type");
673 worker_per_queue = false;
674 }
675 let executor_kind = disk_option.async_executor.unwrap_or_default();
676 let boot_index = disk_option.bootindex;
677 #[cfg(windows)]
678 let io_concurrency = disk_option.io_concurrency.get();
679
680 if block_size % SECTOR_SIZE as u32 != 0 {
681 error!(
682 "Block size {} is not a multiple of {}.",
683 block_size, SECTOR_SIZE,
684 );
685 return Err(SysError::new(libc::EINVAL));
686 }
687 let disk_size = disk_image.get_len()?;
688 if disk_size % block_size as u64 != 0 {
689 warn!(
690 "Disk size {} is not a multiple of block size {}; \
691 the remainder will not be visible to the guest.",
692 disk_size, block_size,
693 );
694 }
695 let num_queues = num_queues.unwrap_or(DEFAULT_NUM_QUEUES);
696 let multi_queue = match num_queues {
697 0 => panic!("Number of queues cannot be zero for a block device"),
698 1 => false,
699 _ => true,
700 };
701 let q_size = queue_size.unwrap_or(DEFAULT_QUEUE_SIZE);
702 if !q_size.is_power_of_two() {
703 error!("queue size {} is not a power of 2.", q_size);
704 return Err(SysError::new(libc::EINVAL));
705 }
706 let queue_sizes = vec![q_size; num_queues as usize];
707
708 let avail_features =
709 Self::build_avail_features(base_features, read_only, sparse, multi_queue, packed_queue);
710
711 let seg_max = get_seg_max(q_size);
712
713 let disk_size = Arc::new(AtomicU64::new(disk_size));
714 let shared_state = Arc::new(AsyncRwLock::new(WorkerSharedState {
715 disk_size: disk_size.clone(),
716 }));
717
718 let mut dontcache_read = disk_option.dontcache;
719 let mut dontcache_write = disk_option.dontcache;
720 if dontcache_write {
721 let descriptors = disk_image.as_raw_descriptors();
722 if descriptors.is_empty()
723 || !descriptors
724 .iter()
725 .all(|&fd| check_dontcache_support(fd, true ))
726 {
727 base::info!(
728 "dontcache requested, but not supported for writes by backing filesystem; falling
729 back to cached I/O"
730 );
731 dontcache_write = false;
732 }
733 }
734 if dontcache_read {
735 let descriptors = disk_image.as_raw_descriptors();
736 if descriptors.is_empty()
737 || !descriptors
738 .iter()
739 .all(|&fd| check_dontcache_support(fd, false ))
740 {
741 base::info!(
742 "dontcache requested, but not supported for reads by backing filesystem; falling
743 back to cached I/O"
744 );
745 dontcache_read = false;
746 }
747 }
748
749 Ok(BlockAsync {
750 disk_image: Some(disk_image),
751 disk_size,
752 avail_features,
753 read_only,
754 sparse,
755 seg_max,
756 block_size,
757 id,
758 queue_sizes,
759 worker_threads: BTreeMap::new(),
760 shared_state,
761 worker_per_queue,
762 control_tube,
763 executor_kind,
764 activated_queues: BTreeSet::new(),
765 boot_index,
766 #[cfg(windows)]
767 io_concurrency,
768 pci_address: disk_option.pci_address,
769 dontcache_read,
770 dontcache_write,
771 })
772 }
773
774 fn build_avail_features(
776 base_features: u64,
777 read_only: bool,
778 sparse: bool,
779 multi_queue: bool,
780 packed_queue: bool,
781 ) -> u64 {
782 let mut avail_features = base_features;
783 if read_only {
784 avail_features |= 1 << VIRTIO_BLK_F_RO;
785 } else {
786 if sparse {
787 avail_features |= 1 << VIRTIO_BLK_F_DISCARD;
788 }
789 avail_features |= 1 << VIRTIO_BLK_F_FLUSH;
790 avail_features |= 1 << VIRTIO_BLK_F_WRITE_ZEROES;
791 }
792 avail_features |= 1 << VIRTIO_BLK_F_SEG_MAX;
793 avail_features |= 1 << VIRTIO_BLK_F_BLK_SIZE;
794 if multi_queue {
795 avail_features |= 1 << VIRTIO_BLK_F_MQ;
796 }
797 if packed_queue {
798 avail_features |= 1 << VIRTIO_F_RING_PACKED;
799 }
800 avail_features
801 }
802
803 async fn execute_request(
808 reader: &mut Reader,
809 writer: &mut Writer,
810 disk_state: &AsyncRwLock<DiskState>,
811 flush_timer: &RefCell<TimerAsync<Timer>>,
812 flush_timer_armed: &RefCell<bool>,
813 ) -> result::Result<(), ExecuteError> {
814 let disk_state = disk_state.read_lock().await;
816 let worker_shared_state = disk_state.worker_shared_state.read_lock().await;
818
819 let req_header: virtio_blk_req_header = reader.read_obj().map_err(ExecuteError::Read)?;
820
821 let req_type = req_header.req_type.to_native();
822 let sector = req_header.sector.to_native();
823
824 if disk_state.read_only && req_type != VIRTIO_BLK_T_IN && req_type != VIRTIO_BLK_T_GET_ID {
825 return Err(ExecuteError::ReadOnly {
826 request_type: req_type,
827 });
828 }
829
830 fn check_range(
833 io_start: u64,
834 io_length: u64,
835 disk_size: u64,
836 ) -> result::Result<(), ExecuteError> {
837 let io_end = io_start
838 .checked_add(io_length)
839 .ok_or(ExecuteError::OutOfRange)?;
840 if io_end > disk_size {
841 Err(ExecuteError::OutOfRange)
842 } else {
843 Ok(())
844 }
845 }
846
847 let disk_size = worker_shared_state.disk_size.load(Ordering::Relaxed);
848 match req_type {
849 VIRTIO_BLK_T_IN => {
850 let data_len = writer.available_bytes();
851 if data_len == 0 {
852 return Ok(());
853 }
854 let offset = sector
855 .checked_shl(u32::from(SECTOR_SHIFT))
856 .ok_or(ExecuteError::OutOfRange)?;
857 check_range(offset, data_len as u64, disk_size)?;
858 let disk_image = &disk_state.disk_image;
859 writer
860 .write_all_from_at_fut(
861 &**disk_image,
862 data_len,
863 offset,
864 IoOptions {
865 dontcache: disk_state.dontcache_read,
866 },
867 )
868 .await
869 .map_err(|desc_error| ExecuteError::ReadIo {
870 length: data_len,
871 sector,
872 desc_error,
873 })?;
874 }
875 VIRTIO_BLK_T_OUT => {
876 let data_len = reader.available_bytes();
877 if data_len == 0 {
878 return Ok(());
879 }
880 let offset = sector
881 .checked_shl(u32::from(SECTOR_SHIFT))
882 .ok_or(ExecuteError::OutOfRange)?;
883 check_range(offset, data_len as u64, disk_size)?;
884 let disk_image = &disk_state.disk_image;
885 reader
886 .read_exact_to_at_fut(
887 &**disk_image,
888 data_len,
889 offset,
890 IoOptions {
891 dontcache: disk_state.dontcache_write,
892 },
893 )
894 .await
895 .map_err(|desc_error| ExecuteError::WriteIo {
896 length: data_len,
897 sector,
898 desc_error,
899 })?;
900
901 if !*flush_timer_armed.borrow() {
902 *flush_timer_armed.borrow_mut() = true;
903
904 let flush_delay = Duration::from_secs(60);
905 flush_timer
906 .borrow_mut()
907 .reset_oneshot(flush_delay)
908 .map_err(ExecuteError::TimerReset)?;
909 }
910 }
911 VIRTIO_BLK_T_DISCARD | VIRTIO_BLK_T_WRITE_ZEROES => {
912 if req_type == VIRTIO_BLK_T_DISCARD && !disk_state.sparse {
913 return Ok(());
915 }
916
917 let seg_count =
918 reader.available_bytes() / size_of::<virtio_blk_discard_write_zeroes>();
919 let seg_max = if req_type == VIRTIO_BLK_T_DISCARD {
920 MAX_DISCARD_SEG as usize
921 } else {
922 MAX_WRITE_ZEROES_SEG as usize
923 };
924 if seg_count > seg_max {
925 return Err(ExecuteError::TooManySegments(seg_count, seg_max));
926 }
927
928 while reader.available_bytes() >= size_of::<virtio_blk_discard_write_zeroes>() {
929 let seg: virtio_blk_discard_write_zeroes =
930 reader.read_obj().map_err(ExecuteError::Read)?;
931
932 let sector = seg.sector.to_native();
933 let num_sectors = seg.num_sectors.to_native();
934 let flags = seg.flags.to_native();
935
936 let valid_flags = if req_type == VIRTIO_BLK_T_WRITE_ZEROES {
937 VIRTIO_BLK_DISCARD_WRITE_ZEROES_FLAG_UNMAP
938 } else {
939 0
940 };
941
942 if (flags & !valid_flags) != 0 {
943 return Err(ExecuteError::DiscardWriteZeroes {
944 ioerr: None,
945 sector,
946 num_sectors,
947 flags,
948 });
949 }
950
951 let offset = sector
952 .checked_shl(u32::from(SECTOR_SHIFT))
953 .ok_or(ExecuteError::OutOfRange)?;
954 let length = u64::from(num_sectors)
955 .checked_shl(u32::from(SECTOR_SHIFT))
956 .ok_or(ExecuteError::OutOfRange)?;
957 check_range(offset, length, disk_size)?;
958
959 if req_type == VIRTIO_BLK_T_DISCARD {
960 let _ = disk_state.disk_image.punch_hole(offset, length).await;
963 } else {
964 disk_state
965 .disk_image
966 .write_zeroes_at(offset, length)
967 .await
968 .map_err(|e| ExecuteError::DiscardWriteZeroes {
969 ioerr: Some(e),
970 sector,
971 num_sectors,
972 flags,
973 })?;
974 }
975 }
976 }
977 VIRTIO_BLK_T_FLUSH => {
978 disk_state
979 .disk_image
980 .fdatasync()
981 .await
982 .map_err(ExecuteError::Flush)?;
983
984 if *flush_timer_armed.borrow() {
985 flush_timer
986 .borrow_mut()
987 .clear()
988 .map_err(ExecuteError::TimerReset)?;
989 *flush_timer_armed.borrow_mut() = false;
990 }
991 }
992 VIRTIO_BLK_T_GET_ID => {
993 writer
994 .write_all(&disk_state.id)
995 .map_err(ExecuteError::CopyId)?;
996 }
997 t => return Err(ExecuteError::Unsupported(t)),
998 };
999 Ok(())
1000 }
1001
1002 fn build_config_space(
1004 disk_size: u64,
1005 seg_max: u32,
1006 block_size: u32,
1007 num_queues: u16,
1008 ) -> virtio_blk_config {
1009 virtio_blk_config {
1010 capacity: Le64::from(disk_size >> SECTOR_SHIFT),
1012 seg_max: Le32::from(seg_max),
1013 blk_size: Le32::from(block_size),
1014 num_queues: Le16::from(num_queues),
1015 max_discard_sectors: Le32::from(MAX_DISCARD_SECTORS),
1016 discard_sector_alignment: Le32::from(DISCARD_SECTOR_ALIGNMENT),
1017 max_write_zeroes_sectors: Le32::from(MAX_WRITE_ZEROES_SECTORS),
1018 write_zeroes_may_unmap: 1,
1019 max_discard_seg: Le32::from(MAX_DISCARD_SEG),
1020 max_write_zeroes_seg: Le32::from(MAX_WRITE_ZEROES_SEG),
1021 ..Default::default()
1022 }
1023 }
1024
1025 #[allow(clippy::map_entry)]
1028 fn start_worker(
1029 &mut self,
1030 idx: usize,
1031 ) -> anyhow::Result<&(WorkerThread<()>, mpsc::UnboundedSender<WorkerCmd>)> {
1032 let key = if self.worker_per_queue { idx } else { 0 };
1033 if self.worker_threads.contains_key(&key) {
1034 return Ok(self.worker_threads.get(&key).unwrap());
1035 }
1036
1037 let ex = self.create_executor();
1038 let control_tube = self.control_tube.take();
1039 let dontcache_read = self.dontcache_read;
1040 let dontcache_write = self.dontcache_write;
1041 let disk_image = if self.worker_per_queue {
1042 self.disk_image
1043 .as_ref()
1044 .context("Failed to ref a disk image")?
1045 .try_clone()
1046 .context("Failed to clone a disk image")?
1047 } else {
1048 self.disk_image
1049 .take()
1050 .context("Failed to take a disk image")?
1051 };
1052 let read_only = self.read_only;
1053 let sparse = self.sparse;
1054 let id = self.id;
1055 let worker_shared_state = self.shared_state.clone();
1056
1057 let (worker_tx, worker_rx) = mpsc::unbounded();
1058 let worker_thread = WorkerThread::start("virtio_blk", move |kill_evt| {
1059 let async_control =
1060 control_tube.map(|c| AsyncTube::new(&ex, c).expect("failed to create async tube"));
1061
1062 let async_image = match disk_image.to_async_disk(&ex) {
1063 Ok(d) => d,
1064 Err(e) => panic!("Failed to create async disk {e:#}"),
1065 };
1066
1067 let disk_state = Rc::new(AsyncRwLock::new(DiskState {
1068 disk_image: async_image,
1069 read_only,
1070 sparse,
1071 id,
1072 dontcache_read,
1073 dontcache_write,
1074 worker_shared_state,
1075 }));
1076
1077 if let Err(err_string) = ex
1078 .run_until(async {
1079 let r = run_worker(&ex, &disk_state, &async_control, worker_rx, kill_evt).await;
1080 if let Err(e) = disk_state.lock().await.disk_image.flush().await {
1082 error!("failed to flush disk image when stopping worker: {e:?}");
1083 }
1084 r
1085 })
1086 .expect("run_until failed")
1087 {
1088 error!("{:#}", err_string);
1089 }
1090 });
1091 match self.worker_threads.entry(key) {
1092 std::collections::btree_map::Entry::Occupied(_) => unreachable!(),
1093 std::collections::btree_map::Entry::Vacant(e) => {
1094 Ok(e.insert((worker_thread, worker_tx)))
1095 }
1096 }
1097 }
1098
1099 pub fn start_queue(
1100 &mut self,
1101 idx: usize,
1102 queue: Queue,
1103 _mem: GuestMemory,
1104 ) -> anyhow::Result<()> {
1105 let (_, worker_tx) = self.start_worker(idx)?;
1106 worker_tx
1107 .unbounded_send(WorkerCmd::StartQueue { index: idx, queue })
1108 .expect("worker channel closed early");
1109 self.activated_queues.insert(idx);
1110 Ok(())
1111 }
1112
1113 pub fn stop_queue(&mut self, idx: usize) -> anyhow::Result<Queue> {
1114 let (_, worker_tx) = self
1117 .worker_threads
1118 .get(if self.worker_per_queue { &idx } else { &0 })
1119 .context("worker not found")?;
1120 let (response_tx, response_rx) = oneshot::channel();
1121 worker_tx
1122 .unbounded_send(WorkerCmd::StopQueue {
1123 index: idx,
1124 response_tx,
1125 })
1126 .expect("worker channel closed early");
1127 let queue = cros_async::block_on(async {
1128 response_rx
1129 .await
1130 .expect("response_rx closed early")
1131 .context("queue not found")
1132 })?;
1133 self.activated_queues.remove(&idx);
1134 Ok(queue)
1135 }
1136}
1137
1138impl VirtioDevice for BlockAsync {
1139 fn keep_rds(&self) -> Vec<RawDescriptor> {
1140 let mut keep_rds = Vec::new();
1141
1142 if let Some(disk_image) = &self.disk_image {
1143 keep_rds.extend(disk_image.as_raw_descriptors());
1144 }
1145
1146 if let Some(control_tube) = &self.control_tube {
1147 keep_rds.push(control_tube.as_raw_descriptor());
1148 }
1149
1150 keep_rds
1151 }
1152
1153 fn features(&self) -> u64 {
1154 self.avail_features
1155 }
1156
1157 fn device_type(&self) -> DeviceType {
1158 DeviceType::Block
1159 }
1160
1161 fn queue_max_sizes(&self) -> &[u16] {
1162 &self.queue_sizes
1163 }
1164
1165 fn read_config(&self, offset: u64, data: &mut [u8]) {
1166 let config_space = {
1167 let disk_size = self.disk_size.load(Ordering::Acquire);
1168 Self::build_config_space(
1169 disk_size,
1170 self.seg_max,
1171 self.block_size,
1172 self.queue_sizes.len() as u16,
1173 )
1174 };
1175 copy_config(data, 0, config_space.as_bytes(), offset);
1176 }
1177
1178 fn activate(
1179 &mut self,
1180 mem: GuestMemory,
1181 _interrupt: Interrupt,
1182 queues: BTreeMap<usize, Queue>,
1183 ) -> anyhow::Result<()> {
1184 for (i, q) in queues {
1185 self.start_queue(i, q, mem.clone())?;
1186 }
1187 Ok(())
1188 }
1189
1190 fn reset(&mut self) -> anyhow::Result<()> {
1191 for (_, (_, worker_tx)) in self.worker_threads.iter_mut() {
1192 let (response_tx, response_rx) = oneshot::channel();
1193 worker_tx
1194 .unbounded_send(WorkerCmd::AbortQueues { response_tx })
1195 .expect("worker channel closed early");
1196 cros_async::block_on(async { response_rx.await.expect("response_rx closed early") });
1197 }
1198 self.activated_queues.clear();
1199 Ok(())
1200 }
1201
1202 fn virtio_sleep(&mut self) -> anyhow::Result<Option<BTreeMap<usize, Queue>>> {
1203 let mut queues = BTreeMap::new();
1205 for index in self.activated_queues.clone() {
1206 queues.insert(index, self.stop_queue(index)?);
1207 }
1208 if queues.is_empty() {
1209 return Ok(None); }
1211 if let Some((_, (_, worker_tx))) = self.worker_threads.iter().next() {
1212 let (response_tx, response_rx) = oneshot::channel();
1213 worker_tx
1214 .unbounded_send(WorkerCmd::FlushDisk { response_tx })
1215 .expect("worker channel closed early");
1216 cros_async::block_on(async { response_rx.await.expect("response_rx closed early") });
1217 }
1218 Ok(Some(queues))
1219 }
1220
1221 fn virtio_wake(
1222 &mut self,
1223 queues_state: Option<(GuestMemory, Interrupt, BTreeMap<usize, Queue>)>,
1224 ) -> anyhow::Result<()> {
1225 if let Some((mem, _interrupt, queues)) = queues_state {
1226 for (i, q) in queues {
1227 self.start_queue(i, q, mem.clone())?
1228 }
1229 }
1230 Ok(())
1231 }
1232
1233 fn virtio_snapshot(&mut self) -> anyhow::Result<AnySnapshot> {
1234 AnySnapshot::to_any(())
1237 }
1238
1239 fn virtio_restore(&mut self, data: AnySnapshot) -> anyhow::Result<()> {
1240 let () = AnySnapshot::from_any(data)?;
1241 Ok(())
1242 }
1243
1244 fn pci_address(&self) -> Option<PciAddress> {
1245 self.pci_address
1246 }
1247
1248 fn bootorder_fw_cfg(&self, pci_slot: u8) -> Option<(Vec<u8>, usize)> {
1249 self.boot_index
1250 .map(|s| (format!("scsi@{pci_slot}/disk@0,0").as_bytes().to_vec(), s))
1251 }
1252}
1253
1254#[cfg(test)]
1255mod tests {
1256 use std::fs::File;
1257 use std::mem::size_of_val;
1258 use std::sync::atomic::AtomicU64;
1259
1260 use data_model::Le32;
1261 use data_model::Le64;
1262 #[cfg(any(target_os = "android", target_os = "linux"))]
1263 use devices::suspendable_virtio_tests;
1264 use devices::virtio::base_features;
1265 use devices::virtio::create_descriptor_chain;
1266 use devices::virtio::DescriptorType;
1267 use devices::virtio::QueueConfig;
1268 use disk::SingleFileDisk;
1269 use hypervisor::ProtectionType;
1270 use tempfile::tempfile;
1271 use tempfile::TempDir;
1272 use vm_memory::GuestAddress;
1273
1274 use super::*;
1275
1276 #[test]
1277 fn read_size() {
1278 let f = tempfile().unwrap();
1279 f.set_len(0x1000).unwrap();
1280
1281 let features = base_features(ProtectionType::Unprotected);
1282 let disk_option = DiskOption::default();
1283 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1284 let mut num_sectors = [0u8; 4];
1285 b.read_config(0, &mut num_sectors);
1286 assert_eq!([0x08, 0x00, 0x00, 0x00], num_sectors);
1288 let mut msw_sectors = [0u8; 4];
1289 b.read_config(4, &mut msw_sectors);
1290 assert_eq!([0x00, 0x00, 0x00, 0x00], msw_sectors);
1292 }
1293
1294 #[test]
1295 fn read_block_size() {
1296 let f = tempfile().unwrap();
1297 f.set_len(0x1000).unwrap();
1298
1299 let features = base_features(ProtectionType::Unprotected);
1300 let disk_option = DiskOption {
1301 block_size: 4096,
1302 sparse: false,
1303 ..Default::default()
1304 };
1305 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1306 let mut blk_size = [0u8; 4];
1307 b.read_config(20, &mut blk_size);
1308 assert_eq!([0x00, 0x10, 0x00, 0x00], blk_size);
1310 }
1311
1312 #[test]
1313 fn read_features() {
1314 let tempdir = TempDir::new().unwrap();
1315 let mut path = tempdir.path().to_owned();
1316 path.push("disk_image");
1317
1318 const DEVICE_FEATURE_BITS: u64 = 0xffffff;
1321
1322 {
1324 let f = File::create(&path).unwrap();
1325 let features = base_features(ProtectionType::Unprotected);
1326 let disk_option = DiskOption::default();
1327 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1328 assert_eq!(0x7244, b.features() & DEVICE_FEATURE_BITS);
1332 }
1333
1334 {
1336 let f = File::create(&path).unwrap();
1337 let features = base_features(ProtectionType::Unprotected);
1338 let disk_option = DiskOption {
1339 sparse: false,
1340 ..Default::default()
1341 };
1342 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1343 assert_eq!(0x5244, b.features() & DEVICE_FEATURE_BITS);
1346 }
1347
1348 {
1350 let f = File::create(&path).unwrap();
1351 let features = base_features(ProtectionType::Unprotected);
1352 let disk_option = DiskOption {
1353 read_only: true,
1354 ..Default::default()
1355 };
1356 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1357 assert_eq!(0x1064, b.features() & DEVICE_FEATURE_BITS);
1360 }
1361 }
1362
1363 #[test]
1364 fn check_pci_adress_configurability() {
1365 let f = tempfile().unwrap();
1366
1367 let features = base_features(ProtectionType::Unprotected);
1368 let disk_option = DiskOption {
1369 pci_address: Some(PciAddress {
1370 bus: 0,
1371 dev: 1,
1372 func: 1,
1373 }),
1374 ..Default::default()
1375 };
1376 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1377
1378 assert_eq!(b.pci_address(), disk_option.pci_address);
1379 }
1380
1381 #[test]
1382 fn check_runtime_blk_queue_configurability() {
1383 let tempdir = TempDir::new().unwrap();
1384 let mut path = tempdir.path().to_owned();
1385 path.push("disk_image");
1386 let features = base_features(ProtectionType::Unprotected);
1387
1388 let f = File::create(&path).unwrap();
1390 let disk_option = DiskOption::default();
1391 let b = BlockAsync::new(features, Box::new(f), &disk_option, None, None, None).unwrap();
1392 assert_eq!(
1393 [DEFAULT_QUEUE_SIZE; DEFAULT_NUM_QUEUES as usize],
1394 b.queue_max_sizes()
1395 );
1396
1397 let f = File::create(&path).unwrap();
1399 let disk_option = DiskOption::default();
1400 let b = BlockAsync::new(
1401 features,
1402 Box::new(f),
1403 &disk_option,
1404 None,
1405 Some(128),
1406 Some(1),
1407 )
1408 .unwrap();
1409 assert_eq!([128; 1], b.queue_max_sizes());
1410 assert_eq!(0, b.features() & (1 << VIRTIO_BLK_F_MQ) as u64);
1412 }
1413
1414 #[test]
1415 fn read_last_sector() {
1416 let ex = Executor::new().expect("creating an executor failed");
1417
1418 let f = tempfile().unwrap();
1419 let disk_size = 0x1000;
1420 f.set_len(disk_size).unwrap();
1421 let af = SingleFileDisk::new(f, &ex).expect("Failed to create SFD");
1422
1423 let mem = Rc::new(
1424 GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1425 .expect("Creating guest memory failed."),
1426 );
1427
1428 let req_hdr = virtio_blk_req_header {
1429 req_type: Le32::from(VIRTIO_BLK_T_IN),
1430 reserved: Le32::from(0),
1431 sector: Le64::from(7), };
1433 mem.write_obj_at_addr(req_hdr, GuestAddress(0x1000))
1434 .expect("writing req failed");
1435
1436 let mut avail_desc = create_descriptor_chain(
1437 &mem,
1438 GuestAddress(0x100), GuestAddress(0x1000), vec![
1441 (DescriptorType::Readable, size_of_val(&req_hdr) as u32),
1443 (DescriptorType::Writable, 512),
1445 (DescriptorType::Writable, 1),
1447 ],
1448 0,
1449 )
1450 .expect("create_descriptor_chain failed");
1451
1452 let timer = Timer::new().expect("Failed to create a timer");
1453 let flush_timer = Rc::new(RefCell::new(
1454 TimerAsync::new(timer, &ex).expect("Failed to create an async timer"),
1455 ));
1456 let flush_timer_armed = Rc::new(RefCell::new(false));
1457
1458 let disk_state = Rc::new(AsyncRwLock::new(DiskState {
1459 disk_image: Box::new(af),
1460 read_only: false,
1461 sparse: true,
1462 id: Default::default(),
1463 dontcache_read: false,
1464 dontcache_write: false,
1465 worker_shared_state: Arc::new(AsyncRwLock::new(WorkerSharedState {
1466 disk_size: Arc::new(AtomicU64::new(disk_size)),
1467 })),
1468 }));
1469
1470 let fut = process_one_request(
1471 &mut avail_desc,
1472 &disk_state,
1473 &flush_timer,
1474 &flush_timer_armed,
1475 );
1476
1477 ex.run_until(fut)
1478 .expect("running executor failed")
1479 .expect("execute failed");
1480
1481 let status_offset = GuestAddress((0x1000 + size_of_val(&req_hdr) + 512) as u64);
1482 let status = mem.read_obj_from_addr::<u8>(status_offset).unwrap();
1483 assert_eq!(status, VIRTIO_BLK_S_OK);
1484 }
1485
1486 #[test]
1487 fn read_beyond_last_sector() {
1488 let f = tempfile().unwrap();
1489 let disk_size = 0x1000;
1490 f.set_len(disk_size).unwrap();
1491 let mem = Rc::new(
1492 GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1493 .expect("Creating guest memory failed."),
1494 );
1495
1496 let req_hdr = virtio_blk_req_header {
1497 req_type: Le32::from(VIRTIO_BLK_T_IN),
1498 reserved: Le32::from(0),
1499 sector: Le64::from(7), };
1501 mem.write_obj_at_addr(req_hdr, GuestAddress(0x1000))
1502 .expect("writing req failed");
1503
1504 let mut avail_desc = create_descriptor_chain(
1505 &mem,
1506 GuestAddress(0x100), GuestAddress(0x1000), vec![
1509 (DescriptorType::Readable, size_of_val(&req_hdr) as u32),
1511 (DescriptorType::Writable, 512 * 2),
1513 (DescriptorType::Writable, 1),
1515 ],
1516 0,
1517 )
1518 .expect("create_descriptor_chain failed");
1519
1520 let ex = Executor::new().expect("creating an executor failed");
1521
1522 let af = SingleFileDisk::new(f, &ex).expect("Failed to create SFD");
1523 let timer = Timer::new().expect("Failed to create a timer");
1524 let flush_timer = Rc::new(RefCell::new(
1525 TimerAsync::new(timer, &ex).expect("Failed to create an async timer"),
1526 ));
1527 let flush_timer_armed = Rc::new(RefCell::new(false));
1528 let disk_state = Rc::new(AsyncRwLock::new(DiskState {
1529 disk_image: Box::new(af),
1530 read_only: false,
1531 sparse: true,
1532 id: Default::default(),
1533 dontcache_read: false,
1534 dontcache_write: false,
1535 worker_shared_state: Arc::new(AsyncRwLock::new(WorkerSharedState {
1536 disk_size: Arc::new(AtomicU64::new(disk_size)),
1537 })),
1538 }));
1539
1540 let fut = process_one_request(
1541 &mut avail_desc,
1542 &disk_state,
1543 &flush_timer,
1544 &flush_timer_armed,
1545 );
1546
1547 ex.run_until(fut)
1548 .expect("running executor failed")
1549 .expect("execute failed");
1550
1551 let status_offset = GuestAddress((0x1000 + size_of_val(&req_hdr) + 512 * 2) as u64);
1552 let status = mem.read_obj_from_addr::<u8>(status_offset).unwrap();
1553 assert_eq!(status, VIRTIO_BLK_S_IOERR);
1554 }
1555
1556 #[test]
1557 fn get_id() {
1558 let ex = Executor::new().expect("creating an executor failed");
1559
1560 let f = tempfile().unwrap();
1561 let disk_size = 0x1000;
1562 f.set_len(disk_size).unwrap();
1563
1564 let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1565 .expect("Creating guest memory failed.");
1566
1567 let req_hdr = virtio_blk_req_header {
1568 req_type: Le32::from(VIRTIO_BLK_T_GET_ID),
1569 reserved: Le32::from(0),
1570 sector: Le64::from(0),
1571 };
1572 mem.write_obj_at_addr(req_hdr, GuestAddress(0x1000))
1573 .expect("writing req failed");
1574
1575 let mut avail_desc = create_descriptor_chain(
1576 &mem,
1577 GuestAddress(0x100), GuestAddress(0x1000), vec![
1580 (DescriptorType::Readable, size_of_val(&req_hdr) as u32),
1582 (DescriptorType::Writable, 20),
1584 (DescriptorType::Writable, 1),
1586 ],
1587 0,
1588 )
1589 .expect("create_descriptor_chain failed");
1590
1591 let af = SingleFileDisk::new(f, &ex).expect("Failed to create SFD");
1592 let timer = Timer::new().expect("Failed to create a timer");
1593 let flush_timer = Rc::new(RefCell::new(
1594 TimerAsync::new(timer, &ex).expect("Failed to create an async timer"),
1595 ));
1596 let flush_timer_armed = Rc::new(RefCell::new(false));
1597
1598 let id = b"a20-byteserialnumber";
1599
1600 let disk_state = Rc::new(AsyncRwLock::new(DiskState {
1601 disk_image: Box::new(af),
1602 read_only: false,
1603 sparse: true,
1604 id: *id,
1605 dontcache_read: false,
1606 dontcache_write: false,
1607 worker_shared_state: Arc::new(AsyncRwLock::new(WorkerSharedState {
1608 disk_size: Arc::new(AtomicU64::new(disk_size)),
1609 })),
1610 }));
1611
1612 let fut = process_one_request(
1613 &mut avail_desc,
1614 &disk_state,
1615 &flush_timer,
1616 &flush_timer_armed,
1617 );
1618
1619 ex.run_until(fut)
1620 .expect("running executor failed")
1621 .expect("execute failed");
1622
1623 let status_offset = GuestAddress((0x1000 + size_of_val(&req_hdr) + 512) as u64);
1624 let status = mem.read_obj_from_addr::<u8>(status_offset).unwrap();
1625 assert_eq!(status, VIRTIO_BLK_S_OK);
1626
1627 let id_offset = GuestAddress(0x1000 + size_of_val(&req_hdr) as u64);
1628 let returned_id = mem.read_obj_from_addr::<[u8; 20]>(id_offset).unwrap();
1629 assert_eq!(returned_id, *id);
1630 }
1631
1632 #[test]
1633 fn reset_and_reactivate_single_worker() {
1634 reset_and_reactivate(false, None);
1635 }
1636
1637 #[test]
1638 fn reset_and_reactivate_multiple_workers() {
1639 reset_and_reactivate(true, None);
1640 }
1641
1642 #[test]
1643 #[cfg(windows)]
1644 fn reset_and_reactivate_overrlapped_io() {
1645 reset_and_reactivate(
1646 false,
1647 Some(
1648 cros_async::sys::windows::ExecutorKindSys::Overlapped { concurrency: None }.into(),
1649 ),
1650 );
1651 }
1652
1653 fn reset_and_reactivate(
1654 enables_multiple_workers: bool,
1655 async_executor: Option<cros_async::ExecutorKind>,
1656 ) {
1657 let f = tempfile::NamedTempFile::new().unwrap();
1659 f.as_file().set_len(0x1000).unwrap();
1660 let path: tempfile::TempPath = f.into_temp_path();
1663
1664 let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1666 .expect("Creating guest memory failed.");
1667
1668 let (_control_tube, control_tube_device) = Tube::pair().unwrap();
1672
1673 let features = base_features(ProtectionType::Unprotected);
1675 let id = b"Block serial number\0";
1676 let disk_option = DiskOption {
1677 path: path.to_path_buf(),
1678 read_only: true,
1679 id: Some(*id),
1680 sparse: false,
1681 multiple_workers: enables_multiple_workers,
1682 async_executor,
1683 ..Default::default()
1684 };
1685 let disk_image = disk_option.open().unwrap();
1686 let mut b = BlockAsync::new(
1687 features,
1688 disk_image,
1689 &disk_option,
1690 Some(control_tube_device),
1691 None,
1692 None,
1693 )
1694 .unwrap();
1695
1696 let interrupt = Interrupt::new_for_test();
1697
1698 let mut q0 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1700 q0.set_ready(true);
1701 let q0 = q0
1702 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1703 .expect("QueueConfig::activate");
1704
1705 let mut q1 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1706 q1.set_ready(true);
1707 let q1 = q1
1708 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1709 .expect("QueueConfig::activate");
1710
1711 b.activate(mem.clone(), interrupt, BTreeMap::from([(0, q0), (1, q1)]))
1712 .expect("activate should succeed");
1713 if !enables_multiple_workers {
1715 assert!(
1716 b.disk_image.is_none(),
1717 "BlockAsync should not have a disk image"
1718 );
1719 }
1720 assert!(
1721 b.control_tube.is_none(),
1722 "BlockAsync should not have a control tube"
1723 );
1724 assert_eq!(
1725 b.worker_threads.len(),
1726 if enables_multiple_workers { 2 } else { 1 }
1727 );
1728
1729 assert!(b.reset().is_ok(), "reset should succeed");
1731 if !enables_multiple_workers {
1732 assert!(
1733 b.disk_image.is_none(),
1734 "BlockAsync should not have a disk image"
1735 );
1736 }
1737 assert!(
1738 b.control_tube.is_none(),
1739 "BlockAsync should not have a control tube"
1740 );
1741 assert_eq!(
1742 b.worker_threads.len(),
1743 if enables_multiple_workers { 2 } else { 1 }
1744 );
1745 assert_eq!(b.id, *b"Block serial number\0");
1746
1747 let interrupt = Interrupt::new_for_test();
1749 let mut q0 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1750 q0.set_ready(true);
1751 let q0 = q0
1752 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1753 .expect("QueueConfig::activate");
1754
1755 let mut q1 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1756 q1.set_ready(true);
1757 let q1 = q1
1758 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1759 .expect("QueueConfig::activate");
1760
1761 b.activate(mem, interrupt, BTreeMap::from([(0, q0), (1, q1)]))
1762 .expect("re-activate should succeed");
1763 }
1764
1765 #[test]
1766 fn resize_with_single_worker() {
1767 resize(false);
1768 }
1769
1770 #[test]
1771 fn resize_with_multiple_workers() {
1772 resize(true);
1774 }
1775
1776 fn resize(enables_multiple_workers: bool) {
1777 let original_size = 0x1000;
1779 let resized_size = 0x2000;
1780
1781 let f = tempfile().unwrap();
1783 f.set_len(original_size).unwrap();
1784 let disk_image: Box<dyn DiskFile> = Box::new(f);
1785 assert_eq!(disk_image.get_len().unwrap(), original_size);
1786
1787 let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1789 .expect("Creating guest memory failed.");
1790
1791 let (control_tube, control_tube_device) = Tube::pair().unwrap();
1793
1794 let features = base_features(ProtectionType::Unprotected);
1796 let disk_option = DiskOption {
1797 multiple_workers: enables_multiple_workers,
1798 ..Default::default()
1799 };
1800 let mut b = BlockAsync::new(
1801 features,
1802 disk_image.try_clone().unwrap(),
1803 &disk_option,
1804 Some(control_tube_device),
1805 None,
1806 None,
1807 )
1808 .unwrap();
1809
1810 let interrupt = Interrupt::new_for_test();
1811
1812 let mut q0 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1814 q0.set_ready(true);
1815 let q0 = q0
1816 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1817 .expect("QueueConfig::activate");
1818
1819 let mut q1 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1820 q1.set_ready(true);
1821 let q1 = q1
1822 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1823 .expect("QueueConfig::activate");
1824
1825 b.activate(mem, interrupt.clone(), BTreeMap::from([(0, q0), (1, q1)]))
1826 .expect("activate should succeed");
1827
1828 assert_eq!(
1830 b.disk_size.load(Ordering::Acquire),
1831 original_size,
1832 "disk_size should be the original size first"
1833 );
1834 let mut capacity = [0u8; 8];
1835 b.read_config(0, &mut capacity);
1836 assert_eq!(
1837 capacity,
1838 [0x8, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00],
1840 "read_config should read the original capacity first"
1841 );
1842
1843 control_tube
1845 .send(&DiskControlCommand::Resize {
1846 new_size: resized_size,
1847 })
1848 .unwrap();
1849 assert_eq!(
1850 control_tube.recv::<DiskControlResult>().unwrap(),
1851 DiskControlResult::Ok,
1852 "resize command should succeed"
1853 );
1854 assert_eq!(
1855 b.disk_size.load(Ordering::Acquire),
1856 resized_size,
1857 "disk_size should be resized to the new size"
1858 );
1859 assert_eq!(
1860 disk_image.get_len().unwrap(),
1861 resized_size,
1862 "underlying disk image should be resized to the new size"
1863 );
1864 let mut capacity = [0u8; 8];
1865 b.read_config(0, &mut capacity);
1866 assert_eq!(
1867 capacity,
1868 [0x10, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00],
1870 "read_config should read the resized capacity"
1871 );
1872 interrupt
1874 .get_interrupt_evt()
1875 .wait()
1876 .expect("interrupt should be signaled");
1877
1878 assert_eq!(
1879 interrupt.read_interrupt_status(),
1880 devices::virtio::INTERRUPT_STATUS_CONFIG_CHANGED as u8,
1881 "INTERRUPT_STATUS_CONFIG_CHANGED should be signaled"
1882 );
1883 }
1884
1885 #[test]
1886 fn run_worker_threads() {
1887 let f = tempfile().unwrap();
1889 f.set_len(0x1000).unwrap();
1890 let disk_image: Box<dyn DiskFile> = Box::new(f);
1891
1892 let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1894 .expect("Creating guest memory failed.");
1895
1896 let features = base_features(ProtectionType::Unprotected);
1898 let disk_option = DiskOption::default();
1899 let mut b = BlockAsync::new(
1900 features,
1901 disk_image.try_clone().unwrap(),
1902 &disk_option,
1903 None,
1904 None,
1905 None,
1906 )
1907 .unwrap();
1908
1909 let interrupt = Interrupt::new_for_test();
1911 let mut q0 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1912 q0.set_ready(true);
1913 let q0 = q0
1914 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1915 .expect("QueueConfig::activate");
1916
1917 let mut q1 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1918 q1.set_ready(true);
1919 let q1 = q1
1920 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1921 .expect("QueueConfig::activate");
1922
1923 b.activate(mem.clone(), interrupt, BTreeMap::from([(0, q0), (1, q1)]))
1924 .expect("activate should succeed");
1925
1926 assert_eq!(b.worker_threads.len(), 1, "1 threads should be spawned.");
1927 drop(b);
1928
1929 let features = base_features(ProtectionType::Unprotected);
1931 let disk_option = DiskOption {
1932 read_only: true,
1933 sparse: false,
1934 multiple_workers: true,
1935 ..DiskOption::default()
1936 };
1937 let mut b = BlockAsync::new(features, disk_image, &disk_option, None, None, None).unwrap();
1938
1939 let interrupt = Interrupt::new_for_test();
1941 let mut q0 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1942 q0.set_ready(true);
1943 let q0 = q0
1944 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1945 .expect("QueueConfig::activate");
1946
1947 let mut q1 = QueueConfig::new(DEFAULT_QUEUE_SIZE, 0);
1948 q1.set_ready(true);
1949 let q1 = q1
1950 .activate(&mem, Event::new().unwrap(), interrupt.clone())
1951 .expect("QueueConfig::activate");
1952
1953 b.activate(mem, interrupt, BTreeMap::from([(0, q0), (1, q1)]))
1954 .expect("activate should succeed");
1955
1956 assert_eq!(b.worker_threads.len(), 2, "2 threads should be spawned.");
1957 }
1958
1959 #[cfg(any(target_os = "android", target_os = "linux"))]
1960 struct BlockContext {}
1961
1962 #[cfg(any(target_os = "android", target_os = "linux"))]
1963 fn modify_device(_block_context: &mut BlockContext, b: &mut BlockAsync) {
1964 b.avail_features = !b.avail_features;
1965 }
1966
1967 #[cfg(any(target_os = "android", target_os = "linux"))]
1968 fn create_device() -> (BlockContext, BlockAsync) {
1969 let f = tempfile().unwrap();
1971 f.set_len(0x1000).unwrap();
1972 let disk_image: Box<dyn DiskFile> = Box::new(f);
1973
1974 let features = base_features(ProtectionType::Unprotected);
1976 let id = b"Block serial number\0";
1977 let disk_option = DiskOption {
1978 read_only: true,
1979 id: Some(*id),
1980 sparse: false,
1981 multiple_workers: true,
1982 ..Default::default()
1983 };
1984 (
1985 BlockContext {},
1986 BlockAsync::new(
1987 features,
1988 disk_image.try_clone().unwrap(),
1989 &disk_option,
1990 None,
1991 None,
1992 None,
1993 )
1994 .unwrap(),
1995 )
1996 }
1997
1998 #[cfg(any(target_os = "android", target_os = "linux"))]
1999 suspendable_virtio_tests!(asyncblock, create_device, 2, modify_device);
2000}