device_virtio_block/
asynchronous.rs

1// Copyright 2021 The ChromiumOS Authors
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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;
101// Arbitrary limits for number of discard/write zeroes segments.
102const MAX_DISCARD_SEG: u32 = 32;
103const MAX_WRITE_ZEROES_SEG: u32 = 32;
104// Hard-coded to 64 KiB (in 512-byte sectors) for now,
105// but this should probably be based on cluster size for qcow.
106const 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            // Log disk I/O errors at debug level to avoid flooding the logs.
185            ExecuteError::ReadIo { .. }
186            | ExecuteError::WriteIo { .. }
187            | ExecuteError::Flush { .. }
188            | ExecuteError::DiscardWriteZeroes { .. } => LogLevel::Debug,
189            // Log all other failures as errors.
190            _ => LogLevel::Error,
191        }
192    }
193}
194
195/// Errors that happen in block outside of executing a request.
196/// This includes errors during resize and flush operations.
197#[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
206/// Maximum length of the virtio-block ID string field.
207const ID_LEN: usize = 20;
208
209/// Virtio block device identifier.
210/// This is an ASCII string terminated by a \0, unless all 20 bytes are used,
211/// in which case the \0 terminator is omitted.
212type BlockId = [u8; ID_LEN];
213
214/// Tracks the state of an anynchronous disk.
215struct 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    /// A DiskState is owned by each worker's executor and cannot be shared by workers, thus
223    /// `worker_shared_state` holds the state shared by workers in Arc.
224    worker_shared_state: Arc<AsyncRwLock<WorkerSharedState>>,
225}
226
227/// Disk state which can be modified by other worker threads
228struct 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    // The last byte of the buffer is virtio_blk_req::status.
242    // Split it into a separate Writer so that status_writer is the final byte and
243    // the original writer is left with just the actual block I/O data.
244    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
275/// Process one descriptor chain asynchronously.
276async 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
298// There is one async task running `handle_queue` per virtio queue in use.
299// Receives messages from the guest and queues a task to complete the operations with the async
300// executor.
301async 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        // Wait for the next signal from `evt` and process `background_tasks` in the meantime.
315        //
316        // NOTE: We can't call `evt.next_val()` directly in the `select!` expression. That would
317        // create a new future each time, which, in the completion-based async backends like
318        // io_uring, means we'd submit a new syscall each time (i.e. a race condition on the
319        // eventfd).
320        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                // Process all the descriptors we've already popped from the queue so that we leave
331                // the queue in a consistent state.
332                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    // Acquire exclusive, mutable access to the state so the virtqueue task won't be able to read
385    // the state while resizing.
386    let disk_state = disk_state.lock().await;
387    // Prevent any other worker threads won't be able to do IO.
388    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    // Allocate new space if the disk image is not sparse.
404    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
419/// Periodically flushes the disk when the given timer fires.
420async 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        // Reset armed before calling fdatasync to guarantee that IO requests that started after we
432        // call fdatasync will be committed eventually.
433        *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        // Once the queue is stopped, it will be sent back over `response_tx`.
453        // `None` indicates that there was no queue at the given index.
454        response_tx: oneshot::Sender<Option<Queue>>,
455    },
456    // Stop all queues without recovering the queues' state and without completing any queued up
457    // work .
458    AbortQueues {
459        // Once the queues are stopped, a `()` value will be sent back over `response_tx`.
460        response_tx: oneshot::Sender<()>,
461    },
462    // Flush any in-memory disk image state to file.
463    FlushDisk {
464        // Once the disk is flushed, a `()` value will be sent back over `response_tx`.
465        response_tx: oneshot::Sender<()>,
466    },
467}
468
469// The main worker thread. Initialized the asynchronous worker tasks and passes them to the executor
470// to be processed.
471//
472// `disk_state` is wrapped by `AsyncRwLock`, which provides both shared and exclusive locks. It's
473// because the state can be read from the virtqueue task while the control task is processing a
474// resizing command.
475async 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    // One flush timer per disk.
483    let timer = Timer::new().expect("Failed to create a timer");
484    let flush_timer_armed = Rc::new(RefCell::new(false));
485
486    // Handles control requests.
487    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    // Handle all the queues in one sub-select call.
492    let flush_timer = Rc::new(RefCell::new(
493        TimerAsync::new(
494            // Call try_clone() to share the same underlying FD with the `flush_disk` task.
495            timer.try_clone().expect("Failed to clone flush_timer"),
496            ex,
497        )
498        .expect("Failed to create an async timer"),
499    ));
500
501    // Flushes the disk periodically.
502    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    // Exit if the kill event is triggered.
507    let kill = async_utils::await_and_exit(ex, kill_evt).fuse();
508    pin_mut!(kill);
509
510    // Running queue handlers.
511    let mut queue_handlers = FuturesUnordered::new();
512    // Async stop functions for queue handlers, by queue index.
513    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                            // Ask the handler to stop.
541                            tx.send(()).unwrap_or_else(|_| panic!("queue handler channel closed early"));
542                            // Wait for its return value.
543                            remote_handle
544                        });
545
546                        // If there was already a handler for this index, stop it before adding the
547                        // new handler future.
548                        if let Some(stop_fn) = old_stop_fn {
549                            warn!("Starting new queue handler without stopping old handler");
550                            // Unfortunately we can't just do `stop_fn().await` because the actual
551                            // work we are waiting on is in `queue_handlers`. So, run both.
552                            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                                // NOTE: This await is blocking the select loop. If we want to
567                                // support stopping queues concurrently, then it needs to be moved.
568                                // For now, keep it simple.
569                                //
570                                // Unfortunately we can't just do `stop_fn().await` because the
571                                // actual work we are waiting on is in `queue_handlers`. So, run
572                                // both.
573                                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 this is the last queue, drop references to the interrupt so
582                                // that, when queues are started up again, we'll use the new
583                                // interrupt passed with the first queue.
584                                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
615/// Virtio device for exposing block level read/write operations on a host file.
616pub struct BlockAsync {
617    // We need to make boot_index public bc the field is used by the main crate to determine boot
618    // order
619    boot_index: Option<usize>,
620    // `None` iff `self.worker_per_queue == false` and the worker thread is running.
621    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    // If `worker_per_queue == true`, `worker_threads` contains the worker for each running queue
633    // by index. Otherwise, contains the monolithic worker for all queues at index 0.
634    //
635    // Once a thread is started, we never stop it, except when `BlockAsync` itself is dropped. That
636    // is because we cannot easily convert the `AsyncDisk` back to a `DiskFile` when backed by
637    // Overlapped I/O on Windows because the file becomes permanently associated with the IOCP
638    // instance of the async executor.
639    worker_threads: BTreeMap<usize, (WorkerThread<()>, mpsc::UnboundedSender<WorkerCmd>)>,
640    shared_state: Arc<AsyncRwLock<WorkerSharedState>>,
641    // Whether to run worker threads in parallel for each queue
642    worker_per_queue: bool,
643    // Indices of running queues.
644    // TODO: The worker already tracks this. Only need it here to stop queues on sleep. Maybe add a
645    // worker cmd to stop all at once, then we can delete this field.
646    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    /// Create a new virtio block device that operates on the given AsyncDisk.
656    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        // Automatically disable multiple workers if the disk image can't be cloned.
671        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 /* write */))
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 /* write */))
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    /// Returns the feature flags given the specified attributes.
775    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    // Execute a single block device request.
804    // `writer` includes the data region only; the status byte is not included.
805    // It is up to the caller to convert the result of this function into a status byte
806    // and write it to the expected location in guest memory.
807    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        // Acquire immutable access to prevent tasks from resizing disk.
815        let disk_state = disk_state.read_lock().await;
816        // Acquire immutable access to prevent other worker threads from resizing disk.
817        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        /// Check that a request accesses only data within the disk's current size.
831        /// All parameters are in units of bytes.
832        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                    // Discard is a hint; if this is a non-sparse disk, just ignore it.
914                    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                        // Since Discard is just a hint and some filesystems may not implement
961                        // FALLOC_FL_PUNCH_HOLE, ignore punch_hole errors.
962                        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    /// Builds and returns the config structure used to specify block features.
1003    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            // If the image is not a multiple of the sector size, the tail bits are not exposed.
1011            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    /// Get the worker for a queue, starting it if necessary.
1026    // NOTE: Can't use `BTreeMap::entry` because it requires an exclusive ref for the whole branch.
1027    #[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                    // Flush any in-memory disk image state to file.
1081                    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        // TODO: Consider stopping the worker thread if this is the last queue managed by it. Then,
1115        // simplify `virtio_sleep` and/or `reset` methods.
1116        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        // Reclaim the queues from workers.
1204        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); // Not activated.
1210        }
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        // `virtio_sleep` ensures there is no pending state, except for the `Queue`s, which are
1235        // handled at a higher layer.
1236        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        // size is 0x1000, so num_sectors is 8 (4096/512).
1287        assert_eq!([0x08, 0x00, 0x00, 0x00], num_sectors);
1288        let mut msw_sectors = [0u8; 4];
1289        b.read_config(4, &mut msw_sectors);
1290        // size is 0x1000, so msw_sectors is 0.
1291        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        // blk_size should be 4096 (0x1000).
1309        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        // Feature bits 0-23 and 50-127 are specific for the device type, but
1319        // at the moment crosvm only supports 64 bits of feature bits.
1320        const DEVICE_FEATURE_BITS: u64 = 0xffffff;
1321
1322        // read-write block device
1323        {
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            // writable device should set VIRTIO_BLK_F_FLUSH + VIRTIO_BLK_F_DISCARD
1329            // + VIRTIO_BLK_F_WRITE_ZEROES + VIRTIO_BLK_F_BLK_SIZE + VIRTIO_BLK_F_SEG_MAX
1330            // + VIRTIO_BLK_F_MQ
1331            assert_eq!(0x7244, b.features() & DEVICE_FEATURE_BITS);
1332        }
1333
1334        // read-write block device, non-sparse
1335        {
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            // writable device should set VIRTIO_F_FLUSH + VIRTIO_BLK_F_RO
1344            // + VIRTIO_BLK_F_BLK_SIZE + VIRTIO_BLK_F_SEG_MAX + VIRTIO_BLK_F_MQ
1345            assert_eq!(0x5244, b.features() & DEVICE_FEATURE_BITS);
1346        }
1347
1348        // read-only block device
1349        {
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            // read-only device should set VIRTIO_BLK_F_RO
1358            // + VIRTIO_BLK_F_BLK_SIZE + VIRTIO_BLK_F_SEG_MAX + VIRTIO_BLK_F_MQ
1359            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        // Default case
1389        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        // Single queue of size 128
1398        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        // Single queue device should not set VIRTIO_BLK_F_MQ
1411        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), // Disk is 8 sectors long, so this is the last valid sector.
1432        };
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),  // Place descriptor chain at 0x100.
1439            GuestAddress(0x1000), // Describe buffer at 0x1000.
1440            vec![
1441                // Request header
1442                (DescriptorType::Readable, size_of_val(&req_hdr) as u32),
1443                // I/O buffer (1 sector of data)
1444                (DescriptorType::Writable, 512),
1445                // Request status
1446                (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), // Disk is 8 sectors long, so this is the last valid sector.
1500        };
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),  // Place descriptor chain at 0x100.
1507            GuestAddress(0x1000), // Describe buffer at 0x1000.
1508            vec![
1509                // Request header
1510                (DescriptorType::Readable, size_of_val(&req_hdr) as u32),
1511                // I/O buffer (2 sectors of data - overlap the end of the disk).
1512                (DescriptorType::Writable, 512 * 2),
1513                // Request status
1514                (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),  // Place descriptor chain at 0x100.
1578            GuestAddress(0x1000), // Describe buffer at 0x1000.
1579            vec![
1580                // Request header
1581                (DescriptorType::Readable, size_of_val(&req_hdr) as u32),
1582                // I/O buffer (20 bytes for serial)
1583                (DescriptorType::Writable, 20),
1584                // Request status
1585                (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        // Create an empty disk image
1658        let f = tempfile::NamedTempFile::new().unwrap();
1659        f.as_file().set_len(0x1000).unwrap();
1660        // Close the file so that it is possible for the disk implementation to take exclusive
1661        // access when opening it.
1662        let path: tempfile::TempPath = f.into_temp_path();
1663
1664        // Create an empty guest memory
1665        let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1666            .expect("Creating guest memory failed.");
1667
1668        // Create a control tube.
1669        // NOTE: We don't want to drop the vmm half of the tube. That would cause the worker thread
1670        // will immediately fail, which isn't what we want to test in this case.
1671        let (_control_tube, control_tube_device) = Tube::pair().unwrap();
1672
1673        // Create a BlockAsync to test
1674        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        // activate with queues of an arbitrary size.
1699        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        // assert resources are consumed
1714        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        // reset and assert resources are still not back (should be in the worker thread)
1730        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        // re-activate should succeed
1748        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        // Test resize handled by one worker affect the whole state
1773        resize(true);
1774    }
1775
1776    fn resize(enables_multiple_workers: bool) {
1777        // disk image size constants
1778        let original_size = 0x1000;
1779        let resized_size = 0x2000;
1780
1781        // Create an empty disk image
1782        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        // Create an empty guest memory
1788        let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1789            .expect("Creating guest memory failed.");
1790
1791        // Create a control tube
1792        let (control_tube, control_tube_device) = Tube::pair().unwrap();
1793
1794        // Create a BlockAsync to test
1795        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        // activate with queues of an arbitrary size.
1813        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 the original size first
1829        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            // original_size (0x1000) >> SECTOR_SHIFT (9) = 0x8
1839            [0x8, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00],
1840            "read_config should read the original capacity first"
1841        );
1842
1843        // assert resize works
1844        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            // resized_size (0x2000) >> SECTOR_SHIFT (9) = 0x10
1869            [0x10, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00],
1870            "read_config should read the resized capacity"
1871        );
1872        // Wait until the blk signals the interrupt
1873        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        // Create an empty duplicable disk image
1888        let f = tempfile().unwrap();
1889        f.set_len(0x1000).unwrap();
1890        let disk_image: Box<dyn DiskFile> = Box::new(f);
1891
1892        // Create an empty guest memory
1893        let mem = GuestMemory::new(&[(GuestAddress(0u64), 4 * 1024 * 1024)])
1894            .expect("Creating guest memory failed.");
1895
1896        // Create a BlockAsync to test with single worker thread
1897        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        // activate with queues of an arbitrary size.
1910        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        // Create a BlockAsync to test with multiple worker threads
1930        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        // activate should succeed
1940        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        // Create an empty disk image
1970        let f = tempfile().unwrap();
1971        f.set_len(0x1000).unwrap();
1972        let disk_image: Box<dyn DiskFile> = Box::new(f);
1973
1974        // Create a BlockAsync to test
1975        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}