use crate::block::BlockRead;
use crate::error::Result;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
pub struct CountingDevice {
inner: Arc<dyn BlockRead>,
reads: AtomicU64,
bytes: AtomicU64,
}
impl CountingDevice {
pub fn new(inner: Arc<dyn BlockRead>) -> Self {
CountingDevice {
inner,
reads: AtomicU64::new(0),
bytes: AtomicU64::new(0),
}
}
pub fn reads(&self) -> u64 {
self.reads.load(Ordering::Relaxed)
}
pub fn bytes(&self) -> u64 {
self.bytes.load(Ordering::Relaxed)
}
pub fn reset(&self) {
self.reads.store(0, Ordering::Relaxed);
self.bytes.store(0, Ordering::Relaxed);
}
}
impl BlockRead for CountingDevice {
fn read_at(&self, offset: u64, buf: &mut [u8]) -> Result<()> {
self.reads.fetch_add(1, Ordering::Relaxed);
self.bytes.fetch_add(buf.len() as u64, Ordering::Relaxed);
self.inner.read_at(offset, buf)
}
fn size_bytes(&self) -> u64 {
self.inner.size_bytes()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_device::Bytes;
fn device(len: usize) -> Arc<CountingDevice> {
Arc::new(CountingDevice::new(Arc::new(Bytes::new(vec![7u8; len]))))
}
#[test]
fn it_counts_calls_and_the_bytes_they_asked_for() {
let dev = device(4096);
let mut small = [0u8; 8];
let mut block = [0u8; 512];
dev.read_at(0, &mut small).expect("read");
assert_eq!((dev.reads(), dev.bytes()), (1, 8));
dev.read_at(1024, &mut block).expect("read");
assert_eq!(
(dev.reads(), dev.bytes()),
(2, 520),
"two calls, and the bytes are the sum of both buffers"
);
}
#[test]
fn a_failed_read_still_counts() {
let dev = device(16);
let mut buf = [0u8; 64];
assert!(dev.read_at(0, &mut buf).is_err(), "past the end");
assert_eq!(dev.reads(), 1, "the driver asked, so it counts");
}
#[test]
fn bytes_counts_what_was_asked_for_including_a_failed_read() {
let dev = device(16);
let mut buf = [0u8; 64];
assert!(dev.read_at(0, &mut buf).is_err(), "past the end");
assert_eq!(
dev.bytes(),
64,
"the whole buffer the driver presented, not the 16 bytes \
available and not the 0 bytes delivered"
);
assert_eq!(
buf, [0u8; 64],
"this device refused without copying, so 64 bytes were \
charged for a transfer of none"
);
}
#[test]
fn bytes_is_the_request_even_when_the_device_moved_a_prefix() {
use std::sync::atomic::AtomicU64;
static N: AtomicU64 = AtomicU64::new(0);
let path = std::env::temp_dir().join(format!(
"fs_core_counting_prefix_{}_{}.bin",
std::process::id(),
N.fetch_add(1, Ordering::Relaxed)
));
std::fs::write(&path, [7u8; 16]).expect("write the fixture");
let file = crate::FileDevice::open(&path).expect("open the fixture");
assert_eq!(file.size_bytes(), 16, "the fixture is the size it claims");
let dev = CountingDevice::new(Arc::new(file));
let mut buf = [0u8; 64];
let err = dev.read_at(0, &mut buf).expect_err("past the end");
let (reads, bytes) = (dev.reads(), dev.bytes());
let moved = buf.iter().filter(|b| **b == 7).count();
drop(dev);
std::fs::remove_file(&path).expect("remove the fixture");
match err {
crate::Error::ShortRead { want, got, .. } => {
assert_eq!((want, got), (64, 16), "asked 64, got the 16 there were");
}
other => panic!("expected ShortRead, got {other:?}"),
}
assert_eq!(moved, 16, "the file device really did copy its prefix");
assert_eq!(
(reads, bytes),
(1, 64),
"charged the request, exactly as the in-memory device was"
);
}
#[test]
fn resetting_starts_the_measurement_where_the_work_does() {
let dev = device(4096);
let mut buf = [0u8; 64];
dev.read_at(0, &mut buf).expect("the mount's own reads");
dev.reset();
assert_eq!((dev.reads(), dev.bytes()), (0, 0));
dev.read_at(64, &mut buf).expect("the work being measured");
assert_eq!((dev.reads(), dev.bytes()), (1, 64));
}
}