use std::{
io,
sync::{
Arc,
atomic::{AtomicBool, AtomicI64, AtomicU8, AtomicU64, AtomicUsize, Ordering},
},
};
use wbase::AlignedBuf;
use wdev::{BufferPool, Device, Result as DeviceResult, SegmentedDevice};
pub(crate) struct CountingDevice {
inner: SegmentedDevice,
reads: AtomicUsize,
read_bytes: AtomicU64,
}
impl CountingDevice {
#[inline]
pub(crate) fn new(inner: SegmentedDevice) -> Self {
Self {
inner,
reads: AtomicUsize::new(0),
read_bytes: AtomicU64::new(0),
}
}
#[inline]
pub(crate) fn reads(&self) -> usize {
self.reads.load(Ordering::Relaxed)
}
#[inline]
pub(crate) fn read_bytes(&self) -> u64 {
self.read_bytes.load(Ordering::Relaxed)
}
}
impl Device for CountingDevice {
#[inline]
fn sector_size(&self) -> usize {
self.inner.sector_size()
}
#[inline]
fn segment_size(&self) -> Option<u64> {
self.inner.segment_size()
}
#[inline]
fn direct_io(&self) -> bool {
self.inner.direct_io()
}
#[inline]
fn pool(&self) -> &Arc<BufferPool> {
self.inner.pool()
}
async fn write_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.write_aligned(offset, buf).await
}
async fn read_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.reads.fetch_add(1, Ordering::Relaxed);
self
.read_bytes
.fetch_add(buf.capacity() as u64, Ordering::Relaxed);
self.inner.read_aligned(offset, buf).await
}
async fn read_raw(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.reads.fetch_add(1, Ordering::Relaxed);
self
.read_bytes
.fetch_add(buf.capacity() as u64, Ordering::Relaxed);
self.inner.read_raw(offset, buf).await
}
async fn sync(&self) -> DeviceResult<()> {
self.inner.sync().await
}
async fn truncate_until_segment(&self, segment_id: u32) -> DeviceResult<()> {
self.inner.truncate_until_segment(segment_id).await
}
}
pub(crate) struct FaultDevice {
inner: SegmentedDevice,
mode: AtomicU8,
}
pub(crate) const MODE_NORMAL: u8 = 0;
pub(crate) const MODE_FAIL: u8 = 1;
pub(crate) const MODE_SHORT: u8 = 2;
impl FaultDevice {
#[inline]
pub(crate) fn new(inner: SegmentedDevice) -> Self {
Self {
inner,
mode: AtomicU8::new(MODE_NORMAL),
}
}
#[inline]
pub(crate) fn set_mode(&self, mode: u8) {
self.mode.store(mode, Ordering::Relaxed);
}
}
impl Device for FaultDevice {
#[inline]
fn sector_size(&self) -> usize {
self.inner.sector_size()
}
#[inline]
fn segment_size(&self) -> Option<u64> {
self.inner.segment_size()
}
#[inline]
fn direct_io(&self) -> bool {
self.inner.direct_io()
}
#[inline]
fn pool(&self) -> &Arc<BufferPool> {
self.inner.pool()
}
async fn write_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
match self.mode.load(Ordering::Relaxed) {
MODE_FAIL => (
Err(wdev::Error::ReadOnly {
offset,
len: buf.len(),
}),
buf,
),
MODE_SHORT => (Ok(buf.len() / 2), buf),
_ => self.inner.write_aligned(offset, buf).await,
}
}
async fn read_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.read_aligned(offset, buf).await
}
async fn read_raw(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.read_raw(offset, buf).await
}
async fn sync(&self) -> DeviceResult<()> {
self.inner.sync().await
}
async fn truncate_until_segment(&self, segment_id: u32) -> DeviceResult<()> {
self.inner.truncate_until_segment(segment_id).await
}
}
pub(crate) struct SyncThrowDevice {
inner: SegmentedDevice,
arm_read_failure: AtomicBool,
throw_on_read_ordinal: AtomicI64,
read_ordinal: AtomicI64,
read_failure_injected: AtomicBool,
}
impl SyncThrowDevice {
#[inline]
pub(crate) fn new(inner: SegmentedDevice) -> Self {
Self {
inner,
arm_read_failure: AtomicBool::new(false),
throw_on_read_ordinal: AtomicI64::new(-1),
read_ordinal: AtomicI64::new(0),
read_failure_injected: AtomicBool::new(false),
}
}
#[inline]
pub(crate) fn set_arm_read_failure(&self, armed: bool) {
self.arm_read_failure.store(armed, Ordering::Relaxed);
}
#[inline]
pub(crate) fn set_throw_on_read_ordinal(&self, ordinal: i64) {
self.throw_on_read_ordinal.store(ordinal, Ordering::Relaxed);
}
#[inline]
pub(crate) fn read_failure_injected(&self) -> bool {
self.read_failure_injected.load(Ordering::Relaxed)
}
fn read_gate(&self) -> Option<wdev::Error> {
if self.arm_read_failure.load(Ordering::Relaxed) {
return Some(wdev::Error::Io(io::Error::other(
"Simulated synchronous device read failure",
)));
}
let ordinal = self.throw_on_read_ordinal.load(Ordering::Relaxed);
if ordinal >= 0 {
let n = self.read_ordinal.fetch_add(1, Ordering::Relaxed);
if n == ordinal {
self.read_failure_injected.store(true, Ordering::Relaxed);
return Some(wdev::Error::Io(io::Error::other(format!(
"Simulated synchronous device read failure on read ordinal {ordinal}"
))));
}
}
None
}
}
impl Device for SyncThrowDevice {
#[inline]
fn sector_size(&self) -> usize {
self.inner.sector_size()
}
#[inline]
fn segment_size(&self) -> Option<u64> {
self.inner.segment_size()
}
#[inline]
fn direct_io(&self) -> bool {
self.inner.direct_io()
}
#[inline]
fn pool(&self) -> &Arc<BufferPool> {
self.inner.pool()
}
async fn write_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
self.inner.write_aligned(offset, buf).await
}
async fn read_aligned(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
match self.read_gate() {
Some(e) => (Err(e), buf),
None => self.inner.read_aligned(offset, buf).await,
}
}
async fn read_raw(&self, offset: u64, buf: AlignedBuf) -> (DeviceResult<usize>, AlignedBuf) {
match self.read_gate() {
Some(e) => (Err(e), buf),
None => self.inner.read_raw(offset, buf).await,
}
}
async fn sync(&self) -> DeviceResult<()> {
self.inner.sync().await
}
async fn truncate_until_segment(&self, segment_id: u32) -> DeviceResult<()> {
self.inner.truncate_until_segment(segment_id).await
}
}