#[cfg(unix)]
use std::fs::File as SyncFile;
use std::{
ffi::OsString,
fs::{DirEntry, ReadDir, create_dir_all, metadata, read_dir, remove_file as sync_remove_file},
io::{ErrorKind, Result as IoResult},
path::{Path, PathBuf},
sync::{
Arc,
atomic::{AtomicBool, AtomicI32, AtomicU32, AtomicU64, Ordering},
},
};
use compio::{
buf::{BufResult, IntoInner, IoBuf},
fs::{File, OpenOptions, remove_file},
io::{AsyncReadAt, AsyncWriteAt},
};
use futures_util::future::join_all;
use gxhash::{GxBuildHasher, HashMap};
use itoa::Buffer;
use papaya::HashMap as PapayaMap;
use wram::{
AlignedBuf, BufferPool, DEFAULT_SECTOR_SIZE, MIN_SECTOR_SIZE, current_thread_id,
is_valid_sector_size,
};
use crate::{
chunk::{SegmentChunks, segment_mask, segment_shift, validate_aligned_io},
device::Device,
error::{Error, Result},
sys::MAX_SEGMENT_SIZE,
};
pub type FileMap = PapayaMap<(u64, u32), Arc<File>, GxBuildHasher>;
#[cfg(debug_assertions)]
const SYNC_GUARD_SEGMENTS: usize = 128;
#[cfg(debug_assertions)]
#[inline]
fn empty_dirty_segs() -> [AtomicU64; 2] {
[const { AtomicU64::new(0) }; 2]
}
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,
#[cfg(target_os = "linux")]
direct_io_probed: AtomicBool,
#[cfg(windows)]
pending_removes: PapayaMap<u32, (), GxBuildHasher>,
pub capacity: Option<u64>,
pub pool: Arc<BufferPool>,
pub dir_syncs: AtomicU64,
#[cfg(debug_assertions)]
dirty_segs: [AtomicU64; 2],
}
#[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
}
#[inline]
pub(crate) const fn parse_u32_ascii(bytes: &[u8]) -> Option<u32> {
if bytes.is_empty() {
return None;
}
let mut val: u32 = 0;
let mut i = 0;
while i < bytes.len() {
let b = bytes[i];
if !b.is_ascii_digit() {
return None;
}
let Some(v) = val.checked_mul(10) else {
return None;
};
let Some(res) = v.checked_add((b - b'0') as u32) else {
return None;
};
val = res;
i += 1;
}
Some(val)
}
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 Some(id) = parse_u32_ascii(rest) 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 !is_valid_sector_size(sector_size) {
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 !is_valid_sector_size(sector_size) {
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: PapayaMap::builder()
.hasher(GxBuildHasher::default())
.build(),
start_segment: AtomicU32::new(0),
end_segment: AtomicI32::new(-1),
direct_io: AtomicBool::new(cfg!(target_os = "linux")),
#[cfg(target_os = "linux")]
direct_io_probed: AtomicBool::new(false),
#[cfg(windows)]
pending_removes: PapayaMap::builder()
.hasher(GxBuildHasher::default())
.build(),
#[cfg(debug_assertions)]
dirty_segs: empty_dirty_segs(),
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 seg_id_u64 = offset >> segment_shift(seg_size);
let seg_id = u32::try_from(seg_id_u64).map_err(|_| Error::SegmentExceeded(seg_id_u64))?;
Ok((seg_id, offset & segment_mask(seg_size)))
}
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) => {
self.direct_io_probed.store(true, Ordering::Relaxed);
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) => {
if !self.direct_io_probed.swap(true, Ordering::Relaxed) {
self.direct_io.store(false, Ordering::Relaxed);
self.files.pin().retain(|_, _| false);
log::error!(
"Direct I/O 探测失败({}: {e}),设备一次性定型为常规缓存 I/O,此后 Direct 打开失败将直接上抛",
path.display()
);
} else if !self.direct_io.load(Ordering::Acquire) {
log::debug!(
"并发 Direct I/O 探测竞态败方({e}),按定型后的常规缓存 I/O 打开: {}",
path.display()
);
} else {
return Err(Error::from(e));
}
}
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);
#[cfg(debug_assertions)]
self.debug_clear_segment(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();
#[cfg(debug_assertions)]
{
for word in &self.dirty_segs {
word.store(0, Ordering::Relaxed);
}
}
}
#[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 & segment_mask(seg_size);
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<()> {
#[cfg(windows)]
self.retry_pending_removes().await;
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 >> segment_shift(seg_size));
if new_start > 0 {
self.truncate_until_segment(new_start as u32).await?;
}
Ok(())
}
#[cfg(windows)]
async fn retry_pending_removes(&self) {
let ids: Vec<u32> = self
.pending_removes
.pin()
.iter()
.map(|(&id, _)| id)
.collect();
let mut done = Vec::new();
for id in ids {
match remove_file(self.segment_path(id)).await {
Ok(()) => done.push(id),
Err(e) if e.kind() == ErrorKind::NotFound => done.push(id),
Err(_) => {}
}
}
if !done.is_empty() {
let pin = self.pending_removes.pin();
for id in done {
pin.remove(&id);
}
}
}
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 min_seg = self.start_segment.load(Ordering::Relaxed);
#[cfg(debug_assertions)]
let pending = self.debug_dirty_segments();
let files: HashMap<u32, Arc<File>> = self
.files
.pin()
.iter()
.filter(|&(&(_, sid), _)| sid >= min_seg)
.map(|(&(_, sid), f)| (sid, Arc::clone(f)))
.collect();
if files.is_empty() {
return Ok(());
}
if files.len() == 1
&& let Some(file) = files.values().next()
{
let res = if datasync {
file.sync_data().await
} else {
file.sync_all().await
};
res.map_err(Error::from)?;
#[cfg(debug_assertions)]
self.debug_verify_synced(&pending, &files);
return Ok(());
}
let results = join_all(files.values().map(|file| async move {
if datasync {
file.sync_data().await
} else {
file.sync_all().await
}
}))
.await;
let mut first_err = None;
for res in results {
if let Err(e) = res
&& first_err.is_none()
{
first_err = Some(Error::from(e));
}
}
if let Some(err) = first_err {
return Err(err);
}
#[cfg(debug_assertions)]
self.debug_verify_synced(&pending, &files);
Ok(())
}
#[cfg(debug_assertions)]
fn debug_mark_dirty(&self, segment_id: u32) {
let seg = segment_id as usize;
if seg < SYNC_GUARD_SEGMENTS {
self.dirty_segs[seg / 64].fetch_or(1 << (seg % 64), Ordering::Relaxed);
}
}
#[cfg(debug_assertions)]
fn debug_verify_synced(&self, pending: &[u32], synced: &HashMap<u32, Arc<File>>) {
let start_seg = u64::from(self.start_segment.load(Ordering::Relaxed));
for &seg in pending {
let bit = !(1 << (seg as usize % 64));
let covered = synced.contains_key(&seg) || u64::from(seg) < start_seg;
assert!(
covered,
"sync 契约违约:段 {seg} 在册有写入,但本次全局 sync 未覆盖其句柄且段未被截断背书,脏页无人 fsync"
);
self.dirty_segs[seg as usize / 64].fetch_and(bit, Ordering::Relaxed);
}
}
#[cfg(debug_assertions)]
fn debug_clear_segment(&self, segment_id: u32) {
let seg = segment_id as usize;
if seg < SYNC_GUARD_SEGMENTS {
self.dirty_segs[seg / 64].fetch_and(!(1 << (seg % 64)), Ordering::Relaxed);
}
}
#[cfg(debug_assertions)]
pub fn debug_dirty_segments(&self) -> Vec<u32> {
(0..SYNC_GUARD_SEGMENTS)
.filter(|&seg| self.dirty_segs[seg / 64].load(Ordering::Relaxed) & (1 << (seg % 64)) != 0)
.map(|seg| seg as u32)
.collect()
}
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),
}
}
}
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;
#[cfg(debug_assertions)]
if res.is_ok() {
self.debug_mark_dirty(seg_id);
}
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;
#[cfg(debug_assertions)]
self.debug_mark_dirty(chunk.seg_id);
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(());
}
#[cfg(windows)]
self.retry_pending_removes().await;
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 => {}
#[cfg(windows)]
Err(e) => {
self.pending_removes.pin().insert(id, ());
log::warn!("段 {id} 删除失败({e}),已记入延迟删除队列");
}
#[cfg(not(windows))]
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());
}
}
}
}
#[cfg(test)]
mod tests {
use super::parse_u32_ascii;
const CONST_PARSED: Option<u32> = parse_u32_ascii(b"123");
const _: () = assert!(matches!(CONST_PARSED, Some(123)));
#[test]
fn test_parse_u32_ascii() {
assert_eq!(parse_u32_ascii(b"0"), Some(0));
assert_eq!(parse_u32_ascii(b"1"), Some(1));
assert_eq!(parse_u32_ascii(b"42"), Some(42));
assert_eq!(parse_u32_ascii(b"007"), Some(7));
assert_eq!(parse_u32_ascii(b"4294967295"), Some(u32::MAX));
assert_eq!(parse_u32_ascii(b""), None);
assert_eq!(parse_u32_ascii(b" "), None);
assert_eq!(parse_u32_ascii(b"a"), None);
assert_eq!(parse_u32_ascii(b"12a"), None);
assert_eq!(parse_u32_ascii(b"a12"), None);
assert_eq!(parse_u32_ascii(b"-1"), None);
assert_eq!(parse_u32_ascii(b"+1"), None);
assert_eq!(parse_u32_ascii(b"4294967296"), None);
assert_eq!(parse_u32_ascii(b"99999999999"), None);
}
}