Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
267 changes: 264 additions & 3 deletions src/vmm/src/devices/virtio/block/virtio/worker.rs
Original file line number Diff line number Diff line change
@@ -1,17 +1,24 @@
// Copyright 2026 Amazon.com, Inc. or its affiliates. All Rights Reserved.
// SPDX-License-Identifier: Apache-2.0

use event_manager::{EventOps, Events, MutEventSubscriber, SubscriberOps};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, Sender, channel};
use std::sync::{Arc, Mutex};
use std::thread::{self, JoinHandle};
use vmm_sys_util::epoll::EventSet;
use vmm_sys_util::eventfd::EventFd;

use super::device::BlockResources;
use super::io::{BlockIoError, FileEngine, async_io};
use super::metrics::BlockDeviceMetrics;
use super::{FinishedRequest, IoErr, ProcessingResult, Request, VirtioBlockError};
use crate::EventManager;
use crate::devices::virtio::device::ActiveState;
use crate::devices::virtio::queue::InvalidAvailIdx;
use crate::devices::virtio::transport::VirtioInterruptType;
use crate::logger::{IncMetric, error};
use crate::logger::{IncMetric, error, warn};
use crate::rate_limiter::RateLimiter;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};

/// Runtime state and processing logic for an active block device.
#[derive(Debug)]
Expand All @@ -23,6 +30,41 @@ pub(crate) struct BlockWorker {
pub(crate) metrics: Arc<BlockDeviceMetrics>,
}

/// Worker state and control channel side (recv ctl msg)
#[derive(Debug)]
struct ThreadedWorker {
state: WorkerState,
control_evt: EventFd,
from_vmm: Receiver<ControlMsg>,
}

/// Data-path ownership state of the worker thread.
#[allow(clippy::large_enum_variant)]
#[derive(Debug)]
enum WorkerState {
Parked,
Running(BlockWorker),
Finished,
}

#[allow(clippy::large_enum_variant)]
enum ControlMsg {
Start(BlockWorker),
Finish(FlushMode),
}

/// VMM-side handle for controlling and joining a block worker thread.
// Handle gets connected to VirtioBlock in follow-up commits
#[allow(dead_code)]
#[derive(Debug)]
pub(crate) struct WorkerHandle {
to_worker: Sender<ControlMsg>,
control_evt: EventFd,
join: JoinHandle<()>,
queue_evts: Vec<EventFd>,
}

/// Determines how pending I/O is handled during worker teardown.
pub(crate) enum FlushMode {
Drain,
DrainAndFlush,
Expand Down Expand Up @@ -252,3 +294,222 @@ impl BlockWorker {
Ok(self.resources.disk.nsectors)
}
}

// Handle gets connected to VirtioBlock in follow-up commits
#[allow(dead_code)]
impl WorkerHandle {
/// Spawn a parked block worker thread.
pub(crate) fn spawn(queue_evts: Vec<EventFd>, name: String) -> Result<Self, std::io::Error> {
// handle writes and worker reads the control eventfd
let control_evt = EventFd::new(libc::EFD_NONBLOCK)?;
let handle_evt = control_evt.try_clone()?;
let (to_worker, from_vmm) = channel::<ControlMsg>();

let join = thread::Builder::new().name(name).spawn(move || {
let event_manager =
EventManager::new().expect("Failed to create block worker EventManager");
run_worker_loop(event_manager, control_evt, from_vmm);
})?;

Ok(Self {
to_worker,
control_evt: handle_evt,
join,
queue_evts,
})
}

pub(crate) fn queue_events(&self) -> &[EventFd] {
&self.queue_evts
}

/// Transfer data-path resources to the worker thread and start processing.
pub(crate) fn start(&self, worker: BlockWorker) {
if let Err(err) = self.to_worker.send(ControlMsg::Start(worker)) {
error!("Failed to send block worker start message: {:?}", err);
return;
}

if let Err(err) = self.control_evt.write(1) {
error!("Failed to notify block worker: {:?}", err);
}
}

/// Stop the worker and wait for its thread to exit.
pub(crate) fn finish(self, flush_mode: FlushMode) {
if let Err(err) = self.to_worker.send(ControlMsg::Finish(flush_mode)) {
error!("Block worker receiver already dropped: {:?}", err);
} else if let Err(err) = self.control_evt.write(1) {
error!("Block worker control event is closed: {:?}", err);
}

self.join.join().unwrap_or_else(|err| {
error!("Block worker thread panicked during teardown: {:?}", err);
});
}
}

impl ThreadedWorker {
const PROCESS_QUEUE: u32 = 0;
const PROCESS_ASYNC_COMPLETION: u32 = 1;
const PROCESS_CONTROL: u32 = 2;

fn register_control_event(&self, ops: &mut EventOps) {
if let Err(err) = ops.add(Events::with_data(
&self.control_evt,
Self::PROCESS_CONTROL,
EventSet::IN,
)) {
error!("Failed to register block worker control event: {}", err);
}
}

fn register_runtime_events(resources: &BlockResources, ops: &mut EventOps) {
if let Err(err) = ops.add(Events::with_data(
&resources.queue_evts[0],
Self::PROCESS_QUEUE,
EventSet::IN,
)) {
error!("Failed to register queue event: {}", err);
}
if let FileEngine::Async(ref engine) = resources.disk.file_engine
&& let Err(err) = ops.add(Events::with_data(
engine.completion_evt(),
Self::PROCESS_ASYNC_COMPLETION,
EventSet::IN,
))
{
error!("Failed to register IO engine completion event: {}", err);
}
}

fn unregister_runtime_events(resources: &BlockResources, ops: &mut EventOps) {
if let Err(err) = ops.remove(Events::with_data(
&resources.queue_evts[0],
Self::PROCESS_QUEUE,
EventSet::IN,
)) {
error!("Failed to unregister queue event: {}", err);
}
if let FileEngine::Async(ref engine) = resources.disk.file_engine
&& let Err(err) = ops.remove(Events::with_data(
engine.completion_evt(),
Self::PROCESS_ASYNC_COMPLETION,
EventSet::IN,
))
{
error!("Failed to unregister IO engine completion event: {}", err);
}
}

fn process_control_event(&mut self, ops: &mut EventOps) {
if let Err(err) = self.control_evt.read() {
error!("Failed to consume block worker control event: {:?}", err);
if let WorkerState::Running(worker) = &self.state {
worker.metrics.event_fails.inc();
}
return;
}

while let Ok(msg) = self.from_vmm.try_recv() {
match msg {
ControlMsg::Start(worker) => self.start_worker(worker, ops),
ControlMsg::Finish(flush_mode) => self.finish_worker(flush_mode, ops),
}

if self.is_finished() {
break;
}
}
}

fn start_worker(&mut self, worker: BlockWorker, ops: &mut EventOps) {
if !matches!(self.state, WorkerState::Parked) {
warn!("Start requested while block worker is not parked");
return;
}

Self::register_runtime_events(&worker.resources, ops);
self.state = WorkerState::Running(worker);
}

fn finish_worker(&mut self, flush_mode: FlushMode, ops: &mut EventOps) {
if let WorkerState::Running(mut worker) =
std::mem::replace(&mut self.state, WorkerState::Finished)
{
Self::unregister_runtime_events(&worker.resources, ops);
match flush_mode {
FlushMode::Drain => worker.drain(true),
FlushMode::DrainAndFlush => worker.drain_and_flush(true),
}
worker.resources.is_io_engine_throttled = false;
}
}

fn is_finished(&self) -> bool {
matches!(self.state, WorkerState::Finished)
}
}

fn run_worker_loop(
mut event_manager: EventManager,
control_evt: EventFd,
from_vmm: Receiver<ControlMsg>,
) {
let worker = Arc::new(Mutex::new(ThreadedWorker {
state: WorkerState::Parked,
control_evt,
from_vmm,
}));
let subscriber: Arc<Mutex<dyn MutEventSubscriber>> = worker.clone();
event_manager.add_subscriber(subscriber);

loop {
if let Err(err) = event_manager.run() {
error!("Block worker event loop error: {:?}", err);
}
if worker
.lock()
.expect("Poisoned block worker lock")
.is_finished()
{
break;
}
}
}

impl MutEventSubscriber for ThreadedWorker {
fn process(&mut self, event: Events, ops: &mut EventOps) {
let source = event.data();
let event_set = event.event_set();

if !EventSet::IN.contains(event_set) {
warn!(
"Block worker received unknown event: {:?} from source: {:?}",
event_set, source
);
return;
}

if let WorkerState::Running(worker) = &mut self.state {
match source {
Self::PROCESS_QUEUE => worker.process_queue_event(),
Self::PROCESS_ASYNC_COMPLETION => worker.process_async_completion_event(),
Self::PROCESS_CONTROL => self.process_control_event(ops),
_ => warn!("Block: Spurious event received: {:?}", source),
}
} else {
match source {
Self::PROCESS_CONTROL => self.process_control_event(ops),
_ => warn!(
"Block: The device worker is not yet activated. Spurious event received: {:?}",
source
),
}
}
}

fn init(&mut self, ops: &mut EventOps) {
self.register_control_event(ops);
}
}
Loading