use crate::store::platform::fs::{
DirEntryInfo, FileKind, FileStat, PositionedReadError, RealFs, StagedFile, StoreDirLockGuard,
StoreFile, StoreFs,
};
use crate::store::StoreError;
use std::collections::BTreeMap;
use std::io;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
#[derive(Default, Clone, Copy)]
struct DurableState {
durable_len: u64,
}
struct SimState {
rng: Mutex<fastrand::Rng>,
durable: Mutex<BTreeMap<PathBuf, DurableState>>,
fsync_drop_one_in: u32,
enospc_on_copy: Mutex<EnospcSchedule>,
op_fault: Mutex<OpFaultSchedule>,
#[cfg(test)]
read_fault: Mutex<ReadFaultSchedule>,
}
pub(crate) struct SimFs {
inner: Arc<dyn StoreFs>,
state: Arc<SimState>,
}
#[derive(Default)]
struct EnospcSchedule {
fail_at: Option<u32>,
seen: u32,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum CrashOp {
Rename,
RemoveFile,
PersistTemp,
}
#[derive(Default)]
struct OpFaultSchedule {
target: Option<(CrashOp, u32)>,
seen: u32,
}
#[cfg(test)]
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum ReadFaultKind {
Io,
ShortRead {
bytes_read: usize,
},
}
#[cfg(test)]
#[derive(Default)]
struct ReadFaultSchedule {
target: Option<(u32, ReadFaultKind)>,
seen: u32,
}
impl SimState {
fn fsync_dropped(&self) -> bool {
let mut rng = self
.rng
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let roll = rng.u32(..);
self.fsync_drop_one_in != 0 && roll.is_multiple_of(self.fsync_drop_one_in)
}
fn record_durable(&self, path: &Path, len: u64) {
let mut durable = self
.durable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
durable.entry(path.to_path_buf()).or_default().durable_len = len;
}
fn move_durable(&self, from: &Path, to: &Path) {
let mut durable = self
.durable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match durable.remove(from) {
Some(state) => {
durable.insert(to.to_path_buf(), state);
}
None => {
durable.remove(to);
}
}
}
fn forget_durable(&self, path: &Path) {
self.durable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(path);
}
fn op_fault_strikes(&self, op: CrashOp) -> bool {
let mut sched = self
.op_fault
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some((target, fail_at)) = sched.target else {
return false;
};
if target != op {
return false;
}
sched.seen = sched.seen.saturating_add(1);
sched.seen == fail_at
}
fn injected_op_fault(op: CrashOp) -> io::Error {
io::Error::other(format!("SimFs: injected fault on {op:?}"))
}
#[cfg(test)]
fn read_fault_strikes(&self) -> Option<ReadFaultKind> {
let mut sched = self
.read_fault
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (fail_at, kind) = sched.target?;
sched.seen = sched.seen.saturating_add(1);
(sched.seen == fail_at).then_some(kind)
}
fn enospc_strikes_now(&self) -> bool {
let mut sched = self
.enospc_on_copy
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(fail_at) = sched.fail_at else {
return false;
};
sched.seen = sched.seen.saturating_add(1);
sched.seen == fail_at
}
}
impl SimFs {
pub(crate) fn new(seed: u64, fsync_drop_one_in: u32) -> Self {
Self::layered(seed, fsync_drop_one_in, Arc::new(RealFs))
}
pub(crate) fn layered(seed: u64, fsync_drop_one_in: u32, inner: Arc<dyn StoreFs>) -> Self {
Self {
inner,
state: Arc::new(SimState {
rng: Mutex::new(fastrand::Rng::with_seed(seed)),
durable: Mutex::new(BTreeMap::new()),
fsync_drop_one_in,
enospc_on_copy: Mutex::new(EnospcSchedule::default()),
op_fault: Mutex::new(OpFaultSchedule::default()),
#[cfg(test)]
read_fault: Mutex::new(ReadFaultSchedule::default()),
}),
}
}
#[cfg(test)]
pub(crate) fn with_fault_on(self, op: CrashOp, fail_at: u32) -> Self {
self.arm_fault_on(op, fail_at);
self
}
#[cfg(test)]
pub(crate) fn arm_fault_on(&self, op: CrashOp, fail_at: u32) {
let mut sched = self
.state
.op_fault
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
sched.target = Some((op, fail_at));
sched.seen = 0;
}
#[cfg(test)]
pub(crate) fn with_read_fault_on(self, fail_at: u32, kind: ReadFaultKind) -> Self {
self.arm_read_fault_on(fail_at, kind);
self
}
#[cfg(test)]
pub(crate) fn arm_read_fault_on(&self, fail_at: u32, kind: ReadFaultKind) {
let mut sched = self
.state
.read_fault
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
sched.target = Some((fail_at, kind));
sched.seen = 0;
}
pub(crate) fn with_enospc_on_copy(self, fail_at: u32) -> Self {
{
let mut sched = self
.state
.enospc_on_copy
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
sched.fail_at = Some(fail_at);
sched.seen = 0;
}
self
}
#[cfg(test)]
pub(crate) fn durable_len(&self, path: &Path) -> u64 {
self.state
.durable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.get(path)
.map_or(0, |state| state.durable_len)
}
pub(crate) fn crash(&self) {
let durable = self
.state
.durable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for (path, state) in durable.iter() {
let inner_is_real_file = matches!(
self.inner.open_file(path),
Ok(handle) if handle.as_std_file().is_some()
);
if inner_is_real_file
&& crate::store::platform::fs::truncate_file_to(path, state.durable_len).is_ok()
{
continue;
}
self.crash_truncate_via_seam(path, state.durable_len);
}
}
fn crash_truncate_via_seam(&self, path: &Path, durable_len: u64) {
let Ok(full) = self.inner.read(path) else {
return;
};
let keep = usize::try_from(durable_len)
.unwrap_or(usize::MAX)
.min(full.len());
if keep >= full.len() {
return;
}
let Some(dir) = path.parent() else {
return;
};
let Ok(mut tmp) = self.inner.named_temp_in(dir) else {
return;
};
if tmp.write_all(&full[..keep]).is_err() || tmp.sync_all().is_err() {
return;
}
let Ok(admission) = crate::store::platform::sync::admit_current_parent_dir_sync() else {
return;
};
let _ = tmp.persist(path, admission);
}
fn track_materialized_file(&self, path: &Path) {
let Ok(meta) = self.inner.metadata(path) else {
return;
};
self.state.record_durable(path, meta.len);
}
}
struct SimStoreFile {
inner: Box<dyn StoreFile>,
path: PathBuf,
state: Arc<SimState>,
}
impl StoreFile for SimStoreFile {
fn write_all(&mut self, buf: &[u8]) -> io::Result<()> {
self.inner.write_all(buf)
}
fn sync_data(&mut self) -> io::Result<()> {
if self.state.fsync_dropped() {
return Ok(());
}
self.state.record_durable(&self.path, self.inner.len()?);
Ok(())
}
fn sync_all(&mut self) -> io::Result<()> {
if self.state.fsync_dropped() {
return Ok(());
}
self.state.record_durable(&self.path, self.inner.len()?);
Ok(())
}
fn len(&self) -> io::Result<u64> {
self.inner.len()
}
fn read_at(&mut self, offset: u64, buf: &mut [u8]) -> io::Result<usize> {
self.inner.read_at(offset, buf)
}
fn read_exact_at(&mut self, offset: u64, buf: &mut [u8]) -> Result<(), PositionedReadError> {
#[cfg(test)]
if let Some(kind) = self.state.read_fault_strikes() {
return Err(match kind {
ReadFaultKind::Io => PositionedReadError::Io(io::Error::other(
"SimFs: injected positioned-read fault",
)),
ReadFaultKind::ShortRead { bytes_read } => {
PositionedReadError::ShortRead { bytes_read }
}
});
}
self.inner.read_exact_at(offset, buf)
}
fn as_std_file(&self) -> Option<&std::fs::File> {
self.inner.as_std_file()
}
}
struct SimStagedFile {
inner: Box<dyn StagedFile>,
state: Arc<SimState>,
written: u64,
durable_staged: u64,
}
impl StagedFile for SimStagedFile {
fn write_all(&mut self, buf: &[u8]) -> io::Result<()> {
self.inner.write_all(buf)?;
self.written = self
.written
.saturating_add(u64::try_from(buf.len()).unwrap_or(u64::MAX));
Ok(())
}
fn sync_all(&mut self) -> io::Result<()> {
if self.state.fsync_dropped() {
return Ok(());
}
self.inner.sync_all()?;
self.durable_staged = self.written;
Ok(())
}
fn persist(
self: Box<Self>,
final_path: &Path,
admission: crate::store::platform::sync::ParentDirSyncAdmission,
) -> io::Result<()> {
if self.state.op_fault_strikes(CrashOp::PersistTemp) {
return Err(SimState::injected_op_fault(CrashOp::PersistTemp));
}
self.inner.persist(final_path, admission)?;
self.state.record_durable(final_path, self.durable_staged);
Ok(())
}
}
impl StoreFs for SimFs {
fn read_dir(&self, path: &Path) -> io::Result<Vec<DirEntryInfo>> {
self.inner.read_dir(path)
}
fn create_dir_all(&self, path: &Path) -> io::Result<()> {
self.inner.create_dir_all(path)
}
fn create_new_file(&self, path: &Path) -> Result<Box<dyn StoreFile>, StoreError> {
let inner = self.inner.create_new_file(path)?;
self.state
.durable
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.insert(path.to_path_buf(), DurableState::default());
Ok(Box::new(SimStoreFile {
inner,
path: path.to_path_buf(),
state: Arc::clone(&self.state),
}))
}
fn open_file(&self, path: &Path) -> io::Result<Box<dyn StoreFile>> {
let inner = self.inner.open_file(path)?;
Ok(Box::new(SimStoreFile {
inner,
path: path.to_path_buf(),
state: Arc::clone(&self.state),
}))
}
fn sync_parent_dir(&self, _path: &Path) -> Result<(), StoreError> {
Ok(())
}
fn reject_symlink_leaf(&self, path: &Path, purpose: &str) -> Result<(), StoreError> {
self.inner.reject_symlink_leaf(path, purpose)
}
fn read(&self, path: &Path) -> io::Result<Vec<u8>> {
self.inner.read(path)
}
fn canonicalize(&self, path: &Path) -> io::Result<PathBuf> {
self.inner.canonicalize(path)
}
fn symlink_metadata(&self, path: &Path) -> io::Result<FileStat> {
self.inner.symlink_metadata(path)
}
fn cow_copy_file(
&self,
from: &Path,
to: &Path,
preference: crate::store::CopyPreference,
) -> io::Result<crate::store::platform::fs::CowStrategyUsed> {
if self.state.enospc_strikes_now() {
return Err(io::Error::new(
io::ErrorKind::StorageFull,
"SimFs: injected ENOSPC mid-fork on cow_copy_file",
));
}
let used = self.inner.cow_copy_file(from, to, preference)?;
self.track_materialized_file(to);
Ok(used)
}
fn copy(&self, from: &Path, to: &Path) -> io::Result<u64> {
if self.state.enospc_strikes_now() {
return Err(io::Error::new(
io::ErrorKind::StorageFull,
"SimFs: injected ENOSPC mid-fork on copy",
));
}
let bytes = self.inner.copy(from, to)?;
self.track_materialized_file(to);
Ok(bytes)
}
fn metadata(&self, path: &Path) -> io::Result<FileStat> {
self.inner.metadata(path)
}
fn rename(&self, from: &Path, to: &Path) -> io::Result<()> {
if self.state.op_fault_strikes(CrashOp::Rename) {
return Err(SimState::injected_op_fault(CrashOp::Rename));
}
self.inner.rename(from, to)?;
self.state.move_durable(from, to);
Ok(())
}
fn remove_file(&self, path: &Path) -> io::Result<()> {
let is_probe_scratch =
path.file_name()
.and_then(|leaf| leaf.to_str())
.is_some_and(|leaf| {
leaf.starts_with(crate::store::platform::evidence::MMAP_PROBE_PREFIX)
});
if !is_probe_scratch && self.state.op_fault_strikes(CrashOp::RemoveFile) {
return Err(SimState::injected_op_fault(CrashOp::RemoveFile));
}
self.inner.remove_file(path)?;
self.state.forget_durable(path);
Ok(())
}
fn remove_dir_all(&self, path: &Path) -> io::Result<()> {
for entry in self.read_dir(path)? {
let child = path.join(&entry.name);
if entry.kind == FileKind::Dir {
self.remove_dir_all(&child)?;
} else {
self.remove_file_if_present(&child)?;
}
}
self.inner.remove_dir_all(path)
}
fn named_temp_in(&self, dir: &Path) -> io::Result<Box<dyn StagedFile>> {
let inner = self.inner.named_temp_in(dir)?;
Ok(Box::new(SimStagedFile {
inner,
state: Arc::clone(&self.state),
written: 0,
durable_staged: 0,
}))
}
fn try_lock_store_dir(
&self,
lock_path: &Path,
) -> Result<Option<Box<dyn StoreDirLockGuard>>, StoreError> {
self.inner.try_lock_store_dir(lock_path)
}
}
#[cfg(test)]
#[path = "fs_tests.rs"]
mod tests;