mod pipeline;
use std::borrow::Cow;
use std::io::ErrorKind;
use std::ops::Range;
use std::path::{Path, PathBuf};
use std::sync::{Arc, OnceLock};
use std::{fs, io, slice};
use memmap2::MmapRaw;
use parking_lot::Mutex;
use self::pipeline::MmapReadPipeline;
use super::traits::{UniversalReadFileOps, UniversalReadFs, UniversalWriteFileOps};
use super::*;
use crate::common::ext::aligned_vec::ACow;
use crate::common::generic_consts::AccessPattern;
use crate::common::mmap::{Advice, AdviceSetting, MULTI_MMAP_IS_SUPPORTED, Madviseable as _};
#[derive(Debug, Default, Clone, Copy)]
pub struct MmapFs;
impl UniversalReadFileOps for MmapFs {
type ContextConfig = ();
fn from_context(_: ()) -> UioResult<Self> {
Ok(MmapFs)
}
fn list_files(&self, prefix_path: &Path) -> UioResult<Vec<ListedFile>> {
local_file_ops::local_list_files(prefix_path)
}
fn exists(&self, path: &Path) -> UioResult<bool> {
fs_err::exists(path).map_err(UniversalIoError::from)
}
}
impl UniversalWriteFileOps for MmapFs {
type AppendFile = MmapFile;
fn create(&self, path: &Path, expected_length: usize) -> UioResult<()> {
local_file_ops::local_create(path, expected_length)
}
fn create_dir(&self, path: &Path) -> UioResult<()> {
local_file_ops::local_create_dir(path)
}
fn remove(&self, path: &Path) -> UioResult<()> {
local_file_ops::local_remove(path)
}
fn remove_dir(&self, path: &Path) -> UioResult<()> {
local_file_ops::local_remove_dir(path)
}
fn atomic_save(&self, path: &Path, bytes: &[u8]) -> UioResult<()> {
local_file_ops::local_atomic_save(path, bytes)
}
fn open_append(&self, path: impl AsRef<Path>, options: OpenOptions) -> UioResult<MmapFile> {
MmapFile::open_inner(path, options.for_append())
}
}
impl UniversalReadFs for MmapFs {
type File = MmapFile;
type OpenExtra = ();
fn open(
&self,
path: impl AsRef<Path>,
options: OpenOptions,
_extra: (),
) -> UioResult<MmapFile> {
MmapFile::open_inner(path, options)
}
}
#[derive(Debug)]
pub struct MmapFile {
path: PathBuf,
writeable: bool,
populate: bool,
#[cfg_attr(target_os = "linux", expect(dead_code))]
advice: AdviceSetting,
append_file: Arc<OnceLock<fs_err::File>>,
mmap: Arc<Mutex<MmapRaw>>,
mmap_seq: Option<Arc<Mutex<MmapRaw>>>,
len: usize,
ptr: SendSyncPtr,
ptr_seq: SendSyncPtr,
}
#[derive(Debug, Clone, Copy)]
struct SendSyncPtr(*mut u8);
unsafe impl Send for SendSyncPtr {}
unsafe impl Sync for SendSyncPtr {}
impl MmapFile {
pub(super) fn open_inner(path: impl AsRef<Path>, options: OpenOptions) -> UioResult<Self> {
let OpenOptions {
writeable,
need_sequential,
populate,
advice,
} = options;
let populate = match populate {
Populate::Auto => Self::populate_auto(),
Populate::Partial(_) | Populate::No => false,
Populate::PreferBackground | Populate::Blocking => true,
};
let mmap = open_mmap(path.as_ref(), writeable, populate, advice)?;
let ptr = SendSyncPtr(mmap.as_mut_ptr());
let (mmap_seq, len, ptr_seq) = if need_sequential && *MULTI_MMAP_IS_SUPPORTED {
let mmap_seq = open_mmap(
path.as_ref(),
false,
false,
AdviceSetting::Advice(Advice::Sequential),
)?;
let len = std::cmp::min(mmap.len(), mmap_seq.len());
let ptr_seq = SendSyncPtr(mmap_seq.as_mut_ptr());
(Some(mmap_seq), len, ptr_seq)
} else {
(None, mmap.len(), ptr)
};
let mmap = Self {
path: path.as_ref().into(),
writeable,
populate,
advice,
append_file: Arc::new(OnceLock::new()),
mmap: Arc::new(Mutex::new(mmap)),
mmap_seq: mmap_seq.map(|mmap_seq_| Arc::new(Mutex::new(mmap_seq_))),
len,
ptr,
ptr_seq,
};
Ok(mmap)
}
}
impl UniversalRead for MmapFile {
type Fs = MmapFs;
type ReadPipeline<'a, U>
= MmapReadPipeline<'a, U>
where
Self: 'a,
U: UserData;
fn reopen(&mut self) -> UioResult<()> {
let old_len = self.len as u64;
let new_len = fs_err::File::open(self.path())
.map_err(|err| UniversalIoError::extract_not_found(err, self.path()))?
.metadata()?
.len();
if new_len < old_len {
return Err(UniversalIoError::Io(io::Error::new(
ErrorKind::UnexpectedEof,
format!(
"Reopen encountered a smaller file than expected; old_len: {old_len}, new_len: {new_len}"
),
)));
}
if new_len == old_len {
return Ok(());
}
self.remap_to(new_len as usize, self.populate)
}
fn read_bytes<P: AccessPattern>(
&self,
range: Range<u64>,
_access_pattern: P,
_align: usize,
) -> UioResult<ACow<'_>> {
let mmap = self.as_bytes::<P>();
let bytes = read_bytes(mmap, range)?;
Ok(ACow::Borrowed(bytes))
}
fn read_iter<P: AccessPattern, T: Item, U: UserData>(
&self,
ranges: impl IntoIterator<Item = (U, ReadRange)>,
_access_pattern: P,
) -> UioResult<impl Iterator<Item = UioResult<(U, Cow<'_, [T]>)>>> {
let bytes = self.as_bytes::<P>();
Ok(ranges.into_iter().map(move |(user_data, range)| {
let items = read_bytemuck::<T>(bytes, range)?;
Ok((user_data, Cow::Borrowed(items)))
}))
}
fn read_batch<P: AccessPattern, T: Item, U: UserData, E: From<UniversalIoError>>(
&self,
ranges: impl IntoIterator<Item = (U, ReadRange)>,
_access_pattern: P,
mut callback: impl FnMut(U, &[T]) -> Result<(), E>,
) -> Result<(), E> {
let bytes = self.as_bytes::<P>();
for (user_data, range) in ranges {
let items = read_bytemuck::<T>(bytes, range)?;
callback(user_data, items)?;
}
Ok(())
}
fn len<T>(&self) -> UioResult<u64> {
let len = self.len / size_of::<T>();
Ok(len as u64)
}
fn populate(&self) -> UioResult<()> {
self.mmap.lock().populate();
Ok(())
}
fn populate_auto() -> bool {
false
}
fn clear_ram_cache(&self) -> UioResult<()> {
unsafe {
self.mmap.lock().drop_page_tables(&self.path);
if let Some(mmap_seq) = &self.mmap_seq {
mmap_seq.lock().drop_page_tables(&self.path);
}
}
crate::common::fs::clear_disk_cache(&self.path)?;
Ok(())
}
fn kind() -> UniversalKind {
UniversalKind::Mmap
}
}
impl UniversalWrite for MmapFile {
fn write<T: bytemuck::Pod>(&mut self, byte_offset: ByteOffset, items: &[T]) -> UioResult<()> {
let mmap = self.as_bytes_mut();
write(mmap, byte_offset, items)?;
Ok(())
}
fn write_batch<'a, T: bytemuck::Pod>(
&mut self,
offset_data: impl IntoIterator<Item = (ByteOffset, &'a [T])>,
) -> UioResult<()> {
let mmap = self.as_bytes_mut();
for (byte_offset, items) in offset_data {
write(mmap, byte_offset, items)?;
}
Ok(())
}
}
impl UniversalFlush for MmapFile {
fn flusher(&self) -> Flusher {
let mmap = self.mmap.clone();
let append_file = self.append_file.clone();
let flusher = move || {
{
let mmap = mmap.lock();
if mmap.len() > 0 {
mmap.flush()?;
}
}
if let Some(file) = append_file.get() {
file.sync_data()?;
}
Ok(())
};
Box::new(flusher)
}
}
impl UniversalAppend for MmapFile {
fn append<T: bytemuck::Pod>(&mut self, offset: ByteOffset, data: &[T]) -> UioResult<()> {
let bytes: &[u8] = bytemuck::cast_slice(data);
if bytes.is_empty() {
return Ok(());
}
{
let mut fd = self.append_fd()?;
self.check_append_offset(fd, offset)?;
io::Write::write_all(&mut fd, bytes)?;
}
self.grow_mapping(offset + bytes.len() as u64)?;
Ok(())
}
fn append_batch<'a, T: bytemuck::Pod>(
&mut self,
offset: ByteOffset,
items: impl IntoIterator<Item = &'a [T]>,
) -> UioResult<()> {
let (mut slices, total) = local_file_ops::collect_append_slices(items);
if total == 0 {
return Ok(());
}
{
let fd = self.append_fd()?;
self.check_append_offset(fd, offset)?;
local_file_ops::write_all_vectored(fd, &mut slices)?;
}
self.grow_mapping(offset + total as u64)?;
Ok(())
}
}
impl MmapFile {
pub(crate) fn grow_mapping(&mut self, new_len: u64) -> UioResult<()> {
debug_assert!(new_len as usize >= self.len, "grow_mapping cannot shrink");
if new_len as usize == self.len {
return Ok(());
}
self.remap_to(new_len as usize, false)
}
fn remap_to(&mut self, new_len: usize, populate: bool) -> UioResult<()> {
let mut mmap = self.mmap.lock();
let mut mmap_seq = self.mmap_seq.as_ref().map(|m| m.lock());
cfg_select! {
target_os = "linux" => {
let _ = populate;
let remap_options = memmap2::RemapOptions::new().may_move(true);
unsafe {
mmap.remap(new_len, remap_options)?;
mmap_seq
.as_mut()
.map(|m| m.remap(new_len, remap_options))
.transpose()?;
};
let ptr = SendSyncPtr(mmap.as_mut_ptr());
let ptr_seq = mmap_seq
.as_ref()
.map(|m| SendSyncPtr(m.as_mut_ptr()))
.unwrap_or(ptr);
let len = new_len;
}
_ => {
let _ = new_len; *mmap = open_mmap(self.path.as_ref(), self.writeable, populate, self.advice)?;
let ptr = SendSyncPtr(mmap.as_mut_ptr());
let ptr_seq;
let len;
if let Some(mmap_seq) = mmap_seq.as_mut() {
let mmap_seq_ = open_mmap(
self.path(),
false,
false,
AdviceSetting::Advice(Advice::Sequential),
)?;
**mmap_seq = mmap_seq_;
len = std::cmp::min(mmap.len(), mmap_seq.len());
ptr_seq = SendSyncPtr(mmap_seq.as_mut_ptr());
} else {
len = mmap.len();
ptr_seq = ptr;
}
}
}
self.ptr = ptr;
self.ptr_seq = ptr_seq;
self.len = len;
Ok(())
}
fn check_append_offset(&self, fd: &fs_err::File, offset: ByteOffset) -> UioResult<()> {
let file_len = fd.metadata()?.len();
if file_len != offset {
return Err(UniversalIoError::AppendOffsetConflict {
path: self.path.clone(),
offset,
});
}
Ok(())
}
fn append_fd(&self) -> UioResult<&fs_err::File> {
if !self.writeable {
return Err(UniversalIoError::Io(io::Error::new(
ErrorKind::PermissionDenied,
"append requires a handle opened with writeable=true",
)));
}
if let Some(file) = self.append_file.get() {
return Ok(file);
}
let file = fs_err::OpenOptions::new()
.append(true)
.open(&self.path)
.map_err(|err| UniversalIoError::extract_not_found(err, &self.path))?;
Ok(self.append_file.get_or_init(|| file))
}
}
fn open_mmap(
path: &Path,
write: bool,
populate: bool,
advice: AdviceSetting,
) -> UioResult<MmapRaw> {
#[expect(clippy::disallowed_types)]
let file = fs::OpenOptions::new()
.read(true)
.write(write)
.open(path)
.map_err(|err| UniversalIoError::extract_not_found(err, path))?;
let mmap = if write {
memmap2::MmapOptions::new().map_raw(&file)?
} else {
memmap2::MmapOptions::new().map_raw_read_only(&file)?
};
if populate {
mmap.populate();
}
mmap.madvise(advice.resolve())?;
Ok(mmap)
}
impl MmapFile {
pub fn path(&self) -> &Path {
&self.path
}
pub fn disk_bytes(&self) -> std::io::Result<u64> {
Ok(fs_err::metadata(&self.path)?.len())
}
#[cfg(unix)]
pub fn resident_bytes(&self) -> std::io::Result<u64> {
let len = self.len;
if len == 0 {
return Ok(0);
}
let page_size = crate::common::mmap::advice::page_size()
.ok_or_else(|| std::io::Error::other("failed to determine page size"))?;
let num_pages = len.div_ceil(page_size);
let mut vec = vec![0u8; num_pages];
let ret = unsafe { nix::libc::mincore(self.ptr.0.cast(), len, vec.as_mut_ptr().cast()) };
if ret != 0 {
return Err(std::io::Error::last_os_error());
}
let resident_pages = vec.iter().filter(|&&b| b & 1 != 0).count();
let resident_bytes = (resident_pages * page_size).min(len) as u64;
Ok(resident_bytes)
}
#[cfg(unix)]
pub fn probe_memory_stats(path: impl AsRef<Path>) -> std::io::Result<(u64, u64)> {
let fs = MmapFs;
let file = fs
.open(
path,
OpenOptions {
writeable: false,
need_sequential: false,
populate: Populate::No,
advice: AdviceSetting::Advice(Advice::Normal),
},
(),
)
.map_err(|e| std::io::Error::other(e.to_string()))?;
let disk_bytes = file.disk_bytes()?;
let resident_bytes = file.resident_bytes()?;
Ok((disk_bytes, resident_bytes))
}
}
impl MmapFile {
pub(crate) fn as_bytes<P: AccessPattern>(&self) -> &[u8] {
let ptr = if P::IS_SEQUENTIAL {
self.ptr_seq
} else {
self.ptr
};
unsafe { slice::from_raw_parts(ptr.0, self.len) }
}
fn as_bytes_mut(&mut self) -> &mut [u8] {
unsafe { slice::from_raw_parts_mut(self.ptr.0, self.len) }
}
}
#[inline]
pub(crate) fn read_bytes(bytes: &[u8], range: Range<u64>) -> UioResult<&[u8]> {
bytes
.get(range.start as usize..range.end as usize)
.ok_or_else(|| UniversalIoError::OutOfBounds {
start: range.start,
end: range.end,
elements: bytes.len(),
})
}
#[inline]
pub(crate) fn read_bytemuck<T: Item>(bytes: &[u8], range: ReadRange) -> UioResult<&[T]> {
let ReadRange {
byte_offset,
length: items,
} = range;
let start = byte_offset as usize;
let end = start + size_of::<T>() * items as usize;
let bytes = bytes
.get(start..end)
.ok_or_else(|| UniversalIoError::OutOfBounds {
start: start as _,
end: end as _,
elements: bytes.len() / size_of::<T>(),
})?;
let items = bytemuck::cast_slice(bytes);
Ok(items)
}
#[inline]
fn write<T>(mmap: &mut [u8], byte_offset: ByteOffset, items: &[T]) -> UioResult<()>
where
T: bytemuck::Pod,
{
let start = byte_offset as usize;
let end = start + size_of_val(items);
let mmap_len_bytes = mmap.len();
let mmap = mmap
.get_mut(start..end)
.ok_or_else(|| UniversalIoError::OutOfBounds {
start: start as _,
end: end as _,
elements: mmap_len_bytes / size_of::<T>(),
})?;
let bytes = bytemuck::cast_slice(items);
mmap.copy_from_slice(bytes);
Ok(())
}
#[cfg(test)]
mod tests {
#[cfg(target_os = "linux")]
use super::*;
#[cfg(target_os = "linux")]
#[test]
fn clear_ram_cache_evicts_dual_mapped_file() {
use std::io::Write as _;
use crate::common::generic_consts::{Random, Sequential};
const LEN: u64 = 8 * 1024 * 1024;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("dual.dat");
let fs_type = nix::sys::statfs::statfs(dir.path())
.unwrap()
.filesystem_type();
if fs_type == nix::sys::statfs::TMPFS_MAGIC {
eprintln!("skipping: tempdir is on tmpfs, fadvise cannot evict its pages");
return;
}
let mut file = fs_err::File::create(&path).unwrap();
file.write_all(&vec![7u8; LEN as usize]).unwrap();
file.sync_all().unwrap();
drop(file);
let options = OpenOptions {
writeable: false,
need_sequential: true,
populate: Populate::No,
advice: AdviceSetting::Global,
};
let mmap = MmapFs.open(&path, options, ()).unwrap();
for bytes in [
mmap.read_bytes(0..LEN, Random, 1).unwrap(),
mmap.read_bytes(0..LEN, Sequential, 1).unwrap(),
] {
assert_eq!(
bytes.iter().map(|byte| u64::from(*byte)).sum::<u64>(),
7 * LEN
);
}
let (size, resident) = MmapFile::probe_memory_stats(&path).unwrap();
assert_eq!(size, LEN);
assert!(resident > LEN / 2, "not cached: {resident} of {size}");
mmap.clear_ram_cache().unwrap();
let (_, resident) = MmapFile::probe_memory_stats(&path).unwrap();
assert!(resident < LEN / 10, "not evicted: {resident} of {size}");
}
}