use std::{
ffi::OsString,
fs::{
DirEntry, File as SyncFile, ReadDir, create_dir_all, metadata, read_dir,
remove_file as sync_remove_file,
},
io::{ErrorKind, Result as IoResult},
path::{Path, PathBuf},
str,
sync::{
Arc,
atomic::{AtomicBool, AtomicI32, AtomicU32, AtomicU64, Ordering},
},
};
use compio::{
buf::{BufResult, IntoInner, IoBuf},
fs::{File, OpenOptions, remove_file},
io::{AsyncReadAt, AsyncWriteAt},
};
use itoa::Buffer;
use whasher::{GxPapayaMap, new_papaya_map};
use wram::{AlignedBuf, BufferPool, DEFAULT_SECTOR_SIZE, MIN_SECTOR_SIZE, current_thread_id};
use crate::{
chunk::{SegmentChunks, validate_aligned_io},
device::Device,
error::{Error, Result},
sys::MAX_SEGMENT_SIZE,
};
pub type FileMap = GxPapayaMap<(u64, u32), Arc<File>>;
pub struct SegmentedDevice {
pub base_path: PathBuf,
pub segment_size: Option<u64>,
pub sector_size: usize,
pub read_only: bool,
pub preallocate: bool,
pub delete_on_close: bool,
pub files: FileMap,
pub start_segment: AtomicU32,
pub end_segment: AtomicI32,
pub direct_io: AtomicBool,
pub capacity: Option<u64>,
pub pool: Arc<BufferPool>,
pub dir_syncs: AtomicU64,
}
#[cfg(unix)]
fn sync_dir(parent: &Path) -> bool {
match SyncFile::open(parent) {
Ok(dir) => {
let synced = dir.sync_all().is_ok();
if !synced {
log::warn!(
"新建段文件后 fsync 父目录 {} 失败,崩溃后新段可能不可见",
parent.display()
);
}
synced
}
Err(e) => {
log::warn!("新建段文件后打开父目录 {} 失败: {e}", parent.display());
false
}
}
}
#[cfg(not(unix))]
fn sync_dir(_parent: &Path) -> bool {
false
}
struct SegmentEntries<'a> {
prefix: &'a [u8],
read_dir: ReadDir,
}
impl Iterator for SegmentEntries<'_> {
type Item = IoResult<(u32, DirEntry)>;
fn next(&mut self) -> Option<Self::Item> {
loop {
let entry = match self.read_dir.next()? {
Ok(e) => e,
Err(e) => return Some(Err(e)),
};
let name = entry.file_name();
let name_bytes = name.as_encoded_bytes();
let Some(rest) = name_bytes.strip_prefix(self.prefix) else {
continue;
};
let Some(rest) = rest.strip_prefix(b".") else {
continue;
};
let Ok(s) = str::from_utf8(rest) else {
continue;
};
if !s.as_bytes().first().is_some_and(u8::is_ascii_digit) {
continue;
}
let Ok(id) = s.parse::<u32>() else {
continue;
};
return Some(Ok((id, entry)));
}
}
}
impl SegmentedDevice {
pub fn new(
base_path: impl Into<PathBuf>,
segment_size: Option<u64>,
sector_size: usize,
) -> Result<Self> {
if sector_size < MIN_SECTOR_SIZE || !sector_size.is_power_of_two() {
return Err(Error::InvalidSectorSize {
size: sector_size,
min: MIN_SECTOR_SIZE,
});
}
let pool = BufferPool::new(sector_size)?;
Self::with_pool(base_path, segment_size, sector_size, pool)
}
pub fn with_pool(
base_path: impl Into<PathBuf>,
segment_size: Option<u64>,
sector_size: usize,
pool: Arc<BufferPool>,
) -> Result<Self> {
if sector_size < MIN_SECTOR_SIZE || !sector_size.is_power_of_two() {
return Err(Error::InvalidSectorSize {
size: sector_size,
min: MIN_SECTOR_SIZE,
});
}
if let Some(seg_size) = segment_size
&& (seg_size == 0
|| !seg_size.is_power_of_two()
|| seg_size < sector_size as u64
|| seg_size > MAX_SEGMENT_SIZE)
{
return Err(Error::InvalidSegmentSize(seg_size));
}
let base_path = base_path.into();
if let Some(parent) = base_path.parent()
&& !parent.as_os_str().is_empty()
{
create_dir_all(parent)?;
}
Ok(Self {
base_path,
segment_size,
sector_size,
read_only: false,
preallocate: false,
delete_on_close: false,
files: new_papaya_map(),
start_segment: AtomicU32::new(0),
end_segment: AtomicI32::new(-1),
direct_io: AtomicBool::new(cfg!(target_os = "linux")),
capacity: None,
pool,
dir_syncs: AtomicU64::new(0),
})
}
pub fn set_capacity(&mut self, capacity: Option<u64>) -> Result<()> {
if let (Some(cap), Some(seg_size)) = (capacity, self.segment_size)
&& (cap == 0 || cap % seg_size != 0)
{
return Err(Error::InvalidCapacity { capacity: cap });
}
self.capacity = capacity;
Ok(())
}
#[inline]
pub fn set_read_only(&mut self, read_only: bool) -> &mut Self {
self.read_only = read_only;
self
}
#[inline]
pub fn set_preallocate(&mut self, preallocate: bool) -> &mut Self {
self.preallocate = preallocate;
self
}
#[inline]
pub fn set_delete_on_close(&mut self, delete_on_close: bool) -> &mut Self {
self.delete_on_close = delete_on_close;
self
}
#[inline]
pub fn is_read_only(&self) -> bool {
self.read_only
}
#[inline]
pub fn is_preallocate(&self) -> bool {
self.preallocate
}
#[inline]
pub fn is_delete_on_close(&self) -> bool {
self.delete_on_close
}
#[inline]
pub fn single_file(base_path: impl Into<PathBuf>) -> Result<Self> {
Self::new(base_path, None, DEFAULT_SECTOR_SIZE)
}
#[inline]
pub fn segmented(base_path: impl Into<PathBuf>, segment_size: u64) -> Result<Self> {
Self::new(base_path, Some(segment_size), DEFAULT_SECTOR_SIZE)
}
#[inline]
fn parent_dir(&self) -> &Path {
match self.base_path.parent() {
Some(p) if !p.as_os_str().is_empty() => p,
_ => Path::new("."),
}
}
pub fn segment_path(&self, segment_id: u32) -> PathBuf {
match self.segment_size {
Some(_) => {
let mut itoa_buf = Buffer::new();
let seg_str = itoa_buf.format(segment_id);
let base = self.base_path.as_os_str();
let mut path = OsString::with_capacity(base.len() + 1 + seg_str.len());
path.push(base);
path.push(".");
path.push(seg_str);
PathBuf::from(path)
}
None => self.base_path.clone(),
}
}
fn segment_entries(&self) -> IoResult<Option<SegmentEntries<'_>>> {
let Some(file_name) = self.base_path.file_name() else {
return Ok(None);
};
let read_dir = read_dir(self.parent_dir())?;
Ok(Some(SegmentEntries {
prefix: file_name.as_encoded_bytes(),
read_dir,
}))
}
#[inline]
pub fn get_segment_and_offset(&self, offset: u64) -> Result<(u32, u64)> {
match self.segment_size {
Some(seg_size) => {
let shift = seg_size.trailing_zeros();
let mask = seg_size - 1;
let seg_id_u64 = offset >> shift;
let seg_id = u32::try_from(seg_id_u64).map_err(|_| Error::SegmentExceeded(seg_id_u64))?;
let off_in_seg = offset & mask;
Ok((seg_id, off_in_seg))
}
None => Ok((0, offset)),
}
}
#[inline]
fn open_options(read_only: bool) -> OpenOptions {
let mut opts = OpenOptions::new();
opts.read(true);
if read_only {
opts.write(false).create(false);
} else {
opts.write(true).create(true);
}
opts
}
async fn try_preallocate(file: &File, path: &Path, preallocate: Option<u64>) {
if let Some(sz) = preallocate
&& let Err(e) = file.set_len(sz).await
{
log::warn!("段文件 {} 预分配至 {sz} 字节失败: {e}", path.display());
}
}
async fn open_file(
&self,
path: &Path,
read_only: bool,
preallocate: Option<u64>,
) -> Result<File> {
if !read_only
&& let Some(parent) = path.parent()
&& !parent.as_os_str().is_empty()
{
let _ = create_dir_all(parent);
}
#[cfg(target_os = "linux")]
if self.direct_io.load(Ordering::Relaxed) {
let mut opts = Self::open_options(read_only);
opts.custom_flags(libc::O_DIRECT);
match opts.open(path).await {
Ok(file) => {
log::debug!("成功以 Direct I/O (O_DIRECT) 打开文件: {}", path.display());
if !read_only {
Self::try_preallocate(&file, path, preallocate).await;
}
return Ok(file);
}
Err(e) if matches!(e.kind(), ErrorKind::InvalidInput | ErrorKind::Unsupported) => {
self.direct_io.store(false, Ordering::Relaxed);
self.files.pin().retain(|_, _| false);
log::warn!(
"Direct I/O 打开文件 {} 失败 ({e}),回退到常规缓存 I/O",
path.display()
);
}
Err(e) => return Err(Error::from(e)),
}
}
let file = Self::open_options(read_only).open(path).await?;
log::debug!("成功打开文件: {}", path.display());
if !read_only {
Self::try_preallocate(&file, path, preallocate).await;
}
Ok(file)
}
async fn get_or_open_file(&self, segment_id: u32) -> Result<Arc<File>> {
if segment_id < self.start_segment.load(Ordering::SeqCst) {
return Err(Error::SegmentNotFound(segment_id));
}
let key = (current_thread_id(), segment_id);
if let Some(file) = self.files.pin().get(&key) {
return Ok(Arc::clone(file));
}
let path = self.segment_path(segment_id);
let is_new_segment =
!self.read_only && metadata(&path).is_err_and(|e| e.kind() == ErrorKind::NotFound);
let prealloc = if self.preallocate && !self.read_only {
self.segment_size
} else {
None
};
let file = self.open_file(&path, self.read_only, prealloc).await?;
if is_new_segment && sync_dir(self.parent_dir()) {
self.dir_syncs.fetch_add(1, Ordering::Relaxed);
}
let (entry, truncated) = {
let pin = self.files.pin();
let entry = Arc::clone(pin.get_or_insert(key, Arc::new(file)));
let truncated = segment_id < self.start_segment.load(Ordering::SeqCst);
if truncated {
pin.remove(&key);
}
(entry, truncated)
};
if truncated {
if !self.read_only {
let _ = remove_file(&path).await;
}
return Err(Error::SegmentNotFound(segment_id));
}
Ok(entry)
}
pub fn get_file_size(&self, segment_id: u32) -> Result<u64> {
if segment_id < self.start_segment.load(Ordering::SeqCst) {
return Ok(0);
}
let path = self.segment_path(segment_id);
match metadata(&path) {
Ok(meta) => Ok(meta.len()),
Err(e) if e.kind() == ErrorKind::NotFound => Ok(0),
Err(e) => Err(Error::Io(e)),
}
}
pub async fn remove_segment(&self, segment_id: u32) -> Result<()> {
self.files.pin().retain(|&(_, sid), _| sid != segment_id);
let path = self.segment_path(segment_id);
match remove_file(&path).await {
Ok(()) => Ok(()),
Err(e) if e.kind() == ErrorKind::NotFound => Ok(()),
Err(e) => Err(Error::Io(e)),
}
}
pub fn reset(&self) {
self.files.pin().clear();
}
#[inline]
pub fn is_segment_cached(&self, segment_id: u32) -> bool {
self.files.pin().keys().any(|&(_, sid)| sid == segment_id)
}
#[inline]
pub fn cached_handle_count(&self) -> usize {
self.files.pin().len()
}
#[inline]
pub fn cached_handles_for_segment(&self, segment_id: u32) -> usize {
self
.files
.pin()
.keys()
.filter(|&(_, sid)| *sid == segment_id)
.count()
}
#[inline]
pub fn is_cached_empty(&self) -> bool {
self.files.pin().is_empty()
}
#[inline]
pub fn start_segment(&self) -> u32 {
self.start_segment.load(Ordering::SeqCst)
}
#[inline]
pub fn end_segment(&self) -> Option<u32> {
let v = self.end_segment.load(Ordering::SeqCst);
(v >= 0).then_some(v as u32)
}
#[inline]
pub fn capacity(&self) -> Option<u64> {
self.capacity
}
#[inline]
fn within_single_segment(&self, offset: u64, len: usize) -> bool {
match self.segment_size {
None => true,
Some(seg_size) => {
let off_in_seg = offset & (seg_size - 1);
off_in_seg
.checked_add(len as u64)
.is_some_and(|end| end <= seg_size)
}
}
}
pub fn recover(&self) -> Result<()> {
let Some(seg_size) = self.segment_size else {
return Ok(());
};
let mut segids: Vec<u32> = Vec::new();
if let Some(entries) = self.segment_entries()? {
for item in entries {
let (id, entry) = item?;
match entry.metadata() {
Ok(m) => {
let file_size = m.len();
if file_size > seg_size {
return Err(Error::SegmentSizeMismatch {
segment: id,
file_size,
segment_size: seg_size,
});
}
segids.push(id);
}
Err(e) if e.kind() == ErrorKind::NotFound => {}
Err(e) => return Err(Error::Io(e)),
}
}
}
segids.sort_unstable();
let mut prev: i64 = -1;
let mut recovered_start = 0u32;
for id in segids {
if i64::from(id) != prev + 1 {
recovered_start = id;
} else {
let seg = i32::try_from(id).unwrap_or(i32::MAX);
self.end_segment.fetch_max(seg, Ordering::SeqCst);
}
prev = i64::from(id);
}
self
.start_segment
.fetch_max(recovered_start, Ordering::SeqCst);
Ok(())
}
async fn handle_capacity(&self, segment: u32) -> Result<()> {
let seg = i32::try_from(segment).unwrap_or(i32::MAX);
if self.end_segment.fetch_max(seg, Ordering::SeqCst) >= seg {
return Ok(());
}
let (Some(cap), Some(seg_size)) = (self.capacity, self.segment_size) else {
return Ok(());
};
let new_start = (segment as u64).saturating_sub(cap >> seg_size.trailing_zeros());
if new_start > 0 {
self.truncate_until_segment(new_start as u32).await?;
}
Ok(())
}
pub async fn sync(&self) -> Result<()> {
self.sync_internal(false).await
}
pub async fn sync_data(&self) -> Result<()> {
self.sync_internal(true).await
}
async fn sync_internal(&self, datasync: bool) -> Result<()> {
let tid = current_thread_id();
let min_seg = self.start_segment.load(Ordering::Relaxed);
let files: Vec<Arc<File>> = self
.files
.pin()
.iter()
.filter(|&(&(t, sid), _)| t == tid && sid >= min_seg)
.map(|(_, f)| Arc::clone(f))
.collect();
let mut first_err = None;
for file in files {
let res = if datasync {
file.sync_data().await
} else {
file.sync_all().await
};
if let Err(e) = res
&& first_err.is_none()
{
first_err = Some(Error::from(e));
}
}
match first_err {
Some(e) => Err(e),
None => Ok(()),
}
}
async fn read_impl(
&self,
offset: u64,
mut buf: AlignedBuf,
aligned: bool,
) -> (Result<usize>, AlignedBuf) {
let sector_size = self.sector_size;
let target_len = buf.required_len().min(buf.capacity());
if aligned {
if let Err(e) = validate_aligned_io(offset, target_len, &buf, sector_size) {
return (Err(e), buf);
}
} else if offset.checked_add(target_len as u64).is_none() {
return (
Err(Error::OutOfBounds {
offset,
len: target_len,
}),
buf,
);
}
if target_len == 0 {
return (Ok(0), buf);
}
if self.within_single_segment(offset, target_len) {
let (seg_id, start_off) = match self.get_segment_and_offset(offset) {
Ok(v) => v,
Err(e) => return (Err(e), buf),
};
let file = match self.get_or_open_file(seg_id).await {
Ok(f) => f,
Err(e) => return (Err(e), buf),
};
let slice = buf.slice(0..target_len);
let BufResult(res, slice) = file.read_at(slice, start_off).await;
buf = slice.into_inner();
let bytes_read = match res {
Ok(n) => n,
Err(e) => {
unsafe { buf.set_len_unchecked(0) };
return (Err(Error::from(e)), buf);
}
};
unsafe { buf.set_len_unchecked(bytes_read) };
return (Ok(bytes_read), buf);
}
unsafe { buf.set_len_unchecked(target_len) };
let mut total_read = 0;
let mut first_err = None;
for chunk in SegmentChunks::new(offset, target_len, self.segment_size) {
let chunk = match chunk {
Ok(c) => c,
Err(e) => {
first_err = Some(e);
break;
}
};
let file = match self.get_or_open_file(chunk.seg_id).await {
Ok(f) => f,
Err(e) => {
first_err = Some(e);
break;
}
};
let slice = buf.slice(chunk.buf_pos..chunk.buf_pos + chunk.len);
let BufResult(res, slice) = file.read_at(slice, chunk.off_in_seg).await;
buf = slice.into_inner();
match res {
Ok(n) => {
total_read += n;
if n < chunk.len {
break;
}
}
Err(e) => {
first_err = Some(Error::Io(e));
break;
}
}
}
unsafe { buf.set_len_unchecked(total_read) };
match first_err {
Some(e) => (Err(e), buf),
None => (Ok(total_read), buf),
}
}
}
unsafe impl Send for SegmentedDevice {}
unsafe impl Sync for SegmentedDevice {}
impl Device for SegmentedDevice {
#[inline]
fn sector_size(&self) -> usize {
self.sector_size
}
#[inline]
fn segment_size(&self) -> Option<u64> {
self.segment_size
}
#[inline]
fn direct_io(&self) -> bool {
self.direct_io.load(Ordering::Relaxed)
}
#[inline]
fn recover(&self) -> Result<()> {
SegmentedDevice::recover(self)
}
#[inline]
fn start_segment(&self) -> u32 {
self.start_segment()
}
#[inline]
fn end_segment(&self) -> Option<u32> {
self.end_segment()
}
#[inline]
fn capacity(&self) -> Option<u64> {
self.capacity
}
#[inline]
fn pool(&self) -> &Arc<BufferPool> {
&self.pool
}
async fn write_aligned(&self, offset: u64, mut buf: AlignedBuf) -> (Result<usize>, AlignedBuf) {
let sector_size = self.sector_size;
let total_len = buf.len();
if self.read_only {
return (
Err(Error::ReadOnly {
offset,
len: total_len,
}),
buf,
);
}
if let Some(cap) = self.capacity
&& self.segment_size.is_none()
&& offset.saturating_add(total_len as u64) > cap
{
return (
Err(Error::OutOfBounds {
offset,
len: total_len,
}),
buf,
);
}
if let Err(e) = validate_aligned_io(offset, total_len, &buf, sector_size) {
return (Err(e), buf);
}
if total_len == 0 {
return (Ok(0), buf);
}
if self.within_single_segment(offset, total_len) {
let (seg_id, start_off) = match self.get_segment_and_offset(offset) {
Ok(v) => v,
Err(e) => return (Err(e), buf),
};
if let Err(e) = self.handle_capacity(seg_id).await {
return (Err(e), buf);
}
let file = match self.get_or_open_file(seg_id).await {
Ok(f) => f,
Err(e) => return (Err(e), buf),
};
let mut file_ref = &*file;
let BufResult(res, buf) = file_ref.write_at(buf, start_off).await;
return (res.map_err(Error::from), buf);
}
let mut total_written = 0;
for chunk in SegmentChunks::new(offset, total_len, self.segment_size) {
let chunk = match chunk {
Ok(c) => c,
Err(e) => return (Err(e), buf),
};
if let Err(e) = self.handle_capacity(chunk.seg_id).await {
return (Err(e), buf);
}
let file = match self.get_or_open_file(chunk.seg_id).await {
Ok(f) => f,
Err(e) => return (Err(e), buf),
};
let slice = buf.slice(chunk.buf_pos..chunk.buf_pos + chunk.len);
let mut file_ref = &*file;
let BufResult(res, slice) = file_ref.write_at(slice, chunk.off_in_seg).await;
buf = slice.into_inner();
match res {
Ok(n) => {
total_written += n;
if n < chunk.len {
break;
}
}
Err(e) => return (Err(Error::Io(e)), buf),
}
}
(Ok(total_written), buf)
}
#[inline]
async fn read_aligned(&self, offset: u64, buf: AlignedBuf) -> (Result<usize>, AlignedBuf) {
self.read_impl(offset, buf, true).await
}
#[inline]
async fn read_raw(&self, offset: u64, buf: AlignedBuf) -> (Result<usize>, AlignedBuf) {
self.read_impl(offset, buf, false).await
}
#[inline]
fn sync(&self) -> impl Future<Output = Result<()>> {
SegmentedDevice::sync(self)
}
#[inline]
fn sync_data(&self) -> impl Future<Output = Result<()>> {
SegmentedDevice::sync_data(self)
}
async fn truncate_until_segment(&self, segment_id: u32) -> Result<()> {
if self.segment_size.is_none() {
return Ok(());
}
if self.start_segment.fetch_max(segment_id, Ordering::SeqCst) >= segment_id {
return Ok(());
}
self.files.pin().retain(|&(_, sid), _| sid >= segment_id);
if let Some(entries) = self.segment_entries()? {
for item in entries {
let (id, entry) = item?;
if id >= segment_id {
continue;
}
match remove_file(entry.path()).await {
Ok(()) => {}
Err(e) if e.kind() == ErrorKind::NotFound => {}
Err(e) => return Err(Error::Io(e)),
}
}
}
Ok(())
}
#[inline]
fn get_file_size(&self, segment_id: u32) -> Result<u64> {
SegmentedDevice::get_file_size(self, segment_id)
}
#[inline]
fn remove_segment(&self, segment_id: u32) -> impl Future<Output = Result<()>> {
SegmentedDevice::remove_segment(self, segment_id)
}
#[inline]
fn reset(&self) {
SegmentedDevice::reset(self);
}
}
impl Drop for SegmentedDevice {
fn drop(&mut self) {
if !self.delete_on_close {
return;
}
if self.segment_size.is_none() {
let _ = sync_remove_file(&self.base_path);
return;
}
if let Ok(Some(entries)) = self.segment_entries() {
for item in entries.flatten() {
let _ = sync_remove_file(item.1.path());
}
}
}
}