review disk

This commit is contained in:
weisd
2024-09-27 09:08:44 +08:00
6 changed files with 1589 additions and 634 deletions
+1
View File
@@ -13,6 +13,7 @@
- [x] 远程rpc
- [x] 错误类型判断,程序中判断错误类型,如何统一错误
- [x] 优化xlmeta, 自定义msg数据结构
- [x] appendFile, createFile, readFile, walk_dir sync io
## 基础功能
@@ -1,9 +1,10 @@
// automatically generated by the FlatBuffers compiler, do not modify
// @generated
use core::cmp::Ordering;
use core::mem;
use core::cmp::Ordering;
extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow};
@@ -11,114 +12,112 @@ use self::flatbuffers::{EndianScalar, Follow};
#[allow(unused_imports, dead_code)]
pub mod models {
use core::cmp::Ordering;
use core::mem;
use core::mem;
use core::cmp::Ordering;
extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow};
extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow};
pub enum PingBodyOffset {}
#[derive(Copy, Clone, PartialEq)]
pub enum PingBodyOffset {}
#[derive(Copy, Clone, PartialEq)]
pub struct PingBody<'a> {
pub _tab: flatbuffers::Table<'a>,
pub struct PingBody<'a> {
pub _tab: flatbuffers::Table<'a>,
}
impl<'a> flatbuffers::Follow<'a> for PingBody<'a> {
type Inner = PingBody<'a>;
#[inline]
unsafe fn follow(buf: &'a [u8], loc: usize) -> Self::Inner {
Self { _tab: flatbuffers::Table::new(buf, loc) }
}
}
impl<'a> PingBody<'a> {
pub const VT_PAYLOAD: flatbuffers::VOffsetT = 4;
pub const fn get_fully_qualified_name() -> &'static str {
"models.PingBody"
}
#[inline]
pub unsafe fn init_from_table(table: flatbuffers::Table<'a>) -> Self {
PingBody { _tab: table }
}
#[allow(unused_mut)]
pub fn create<'bldr: 'args, 'args: 'mut_bldr, 'mut_bldr, A: flatbuffers::Allocator + 'bldr>(
_fbb: &'mut_bldr mut flatbuffers::FlatBufferBuilder<'bldr, A>,
args: &'args PingBodyArgs<'args>
) -> flatbuffers::WIPOffset<PingBody<'bldr>> {
let mut builder = PingBodyBuilder::new(_fbb);
if let Some(x) = args.payload { builder.add_payload(x); }
builder.finish()
}
#[inline]
pub fn payload(&self) -> Option<flatbuffers::Vector<'a, u8>> {
// Safety:
// Created from valid Table for this object
// which contains a valid value in this slot
unsafe { self._tab.get::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'a, u8>>>(PingBody::VT_PAYLOAD, None)}
}
}
impl flatbuffers::Verifiable for PingBody<'_> {
#[inline]
fn run_verifier(
v: &mut flatbuffers::Verifier, pos: usize
) -> Result<(), flatbuffers::InvalidFlatbuffer> {
use self::flatbuffers::Verifiable;
v.visit_table(pos)?
.visit_field::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'_, u8>>>("payload", Self::VT_PAYLOAD, false)?
.finish();
Ok(())
}
}
pub struct PingBodyArgs<'a> {
pub payload: Option<flatbuffers::WIPOffset<flatbuffers::Vector<'a, u8>>>,
}
impl<'a> Default for PingBodyArgs<'a> {
#[inline]
fn default() -> Self {
PingBodyArgs {
payload: None,
}
}
}
impl<'a> flatbuffers::Follow<'a> for PingBody<'a> {
type Inner = PingBody<'a>;
#[inline]
unsafe fn follow(buf: &'a [u8], loc: usize) -> Self::Inner {
Self {
_tab: flatbuffers::Table::new(buf, loc),
}
}
pub struct PingBodyBuilder<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> {
fbb_: &'b mut flatbuffers::FlatBufferBuilder<'a, A>,
start_: flatbuffers::WIPOffset<flatbuffers::TableUnfinishedWIPOffset>,
}
impl<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> PingBodyBuilder<'a, 'b, A> {
#[inline]
pub fn add_payload(&mut self, payload: flatbuffers::WIPOffset<flatbuffers::Vector<'b , u8>>) {
self.fbb_.push_slot_always::<flatbuffers::WIPOffset<_>>(PingBody::VT_PAYLOAD, payload);
}
#[inline]
pub fn new(_fbb: &'b mut flatbuffers::FlatBufferBuilder<'a, A>) -> PingBodyBuilder<'a, 'b, A> {
let start = _fbb.start_table();
PingBodyBuilder {
fbb_: _fbb,
start_: start,
}
}
#[inline]
pub fn finish(self) -> flatbuffers::WIPOffset<PingBody<'a>> {
let o = self.fbb_.end_table(self.start_);
flatbuffers::WIPOffset::new(o.value())
}
}
impl<'a> PingBody<'a> {
pub const VT_PAYLOAD: flatbuffers::VOffsetT = 4;
impl core::fmt::Debug for PingBody<'_> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
let mut ds = f.debug_struct("PingBody");
ds.field("payload", &self.payload());
ds.finish()
}
}
} // pub mod models
pub const fn get_fully_qualified_name() -> &'static str {
"models.PingBody"
}
#[inline]
pub unsafe fn init_from_table(table: flatbuffers::Table<'a>) -> Self {
PingBody { _tab: table }
}
#[allow(unused_mut)]
pub fn create<'bldr: 'args, 'args: 'mut_bldr, 'mut_bldr, A: flatbuffers::Allocator + 'bldr>(
_fbb: &'mut_bldr mut flatbuffers::FlatBufferBuilder<'bldr, A>,
args: &'args PingBodyArgs<'args>,
) -> flatbuffers::WIPOffset<PingBody<'bldr>> {
let mut builder = PingBodyBuilder::new(_fbb);
if let Some(x) = args.payload {
builder.add_payload(x);
}
builder.finish()
}
#[inline]
pub fn payload(&self) -> Option<flatbuffers::Vector<'a, u8>> {
// Safety:
// Created from valid Table for this object
// which contains a valid value in this slot
unsafe {
self._tab
.get::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'a, u8>>>(PingBody::VT_PAYLOAD, None)
}
}
}
impl flatbuffers::Verifiable for PingBody<'_> {
#[inline]
fn run_verifier(v: &mut flatbuffers::Verifier, pos: usize) -> Result<(), flatbuffers::InvalidFlatbuffer> {
use self::flatbuffers::Verifiable;
v.visit_table(pos)?
.visit_field::<flatbuffers::ForwardsUOffset<flatbuffers::Vector<'_, u8>>>("payload", Self::VT_PAYLOAD, false)?
.finish();
Ok(())
}
}
pub struct PingBodyArgs<'a> {
pub payload: Option<flatbuffers::WIPOffset<flatbuffers::Vector<'a, u8>>>,
}
impl<'a> Default for PingBodyArgs<'a> {
#[inline]
fn default() -> Self {
PingBodyArgs { payload: None }
}
}
pub struct PingBodyBuilder<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> {
fbb_: &'b mut flatbuffers::FlatBufferBuilder<'a, A>,
start_: flatbuffers::WIPOffset<flatbuffers::TableUnfinishedWIPOffset>,
}
impl<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> PingBodyBuilder<'a, 'b, A> {
#[inline]
pub fn add_payload(&mut self, payload: flatbuffers::WIPOffset<flatbuffers::Vector<'b, u8>>) {
self.fbb_
.push_slot_always::<flatbuffers::WIPOffset<_>>(PingBody::VT_PAYLOAD, payload);
}
#[inline]
pub fn new(_fbb: &'b mut flatbuffers::FlatBufferBuilder<'a, A>) -> PingBodyBuilder<'a, 'b, A> {
let start = _fbb.start_table();
PingBodyBuilder {
fbb_: _fbb,
start_: start,
}
}
#[inline]
pub fn finish(self) -> flatbuffers::WIPOffset<PingBody<'a>> {
let o = self.fbb_.end_table(self.start_);
flatbuffers::WIPOffset::new(o.value())
}
}
impl core::fmt::Debug for PingBody<'_> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
let mut ds = f.debug_struct("PingBody");
ds.field("payload", &self.payload());
ds.finish()
}
}
} // pub mod models
File diff suppressed because it is too large Load Diff
+381 -241
View File
@@ -10,7 +10,7 @@ use crate::disk::error::{
};
use crate::disk::os::check_path_length;
use crate::disk::{LocalFileReader, LocalFileWriter, STORAGE_FORMAT_FILE};
use crate::utils::fs::{lstat, O_RDONLY};
use crate::utils::fs::{lstat, O_APPEND, O_CREATE, O_RDONLY, O_WRONLY};
use crate::utils::path::{has_suffix, SLASH_SEPARATOR};
use crate::{
error::{Error, Result},
@@ -27,7 +27,7 @@ use time::OffsetDateTime;
use tokio::fs::{self, File};
use tokio::io::{AsyncReadExt, AsyncWriteExt, ErrorKind};
use tokio::sync::RwLock;
use tracing::{debug, warn};
use tracing::{error, warn};
use uuid::Uuid;
#[derive(Debug)]
@@ -197,21 +197,6 @@ impl LocalDisk {
// })
// }
pub async fn rename_all(&self, src_data_path: &PathBuf, dst_data_path: &PathBuf, skip: &PathBuf) -> Result<()> {
if !skip.starts_with(src_data_path) {
fs::create_dir_all(dst_data_path.parent().unwrap_or(Path::new("/"))).await?;
}
// debug!(
// "rename_all from \n {:?} \n to \n {:?} \n skip:{:?}",
// &src_data_path, &dst_data_path, &skip
// );
fs::rename(&src_data_path, &dst_data_path).await?;
Ok(())
}
pub async fn move_to_trash(&self, delete_path: &PathBuf, _recursive: bool, _immediate_purge: bool) -> Result<()> {
let trash_path = self.get_object_path(super::RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?;
if let Some(parent) = trash_path.parent() {
@@ -332,6 +317,12 @@ impl LocalDisk {
}
}
async fn read_metadata(&self, file_path: impl AsRef<Path>) -> Result<Vec<u8>> {
// TODO: suport timeout
let (data, _) = self.read_metadata_with_dmtime(file_path.as_ref()).await?;
Ok(data)
}
// FIXME: read_metadata only suport
async fn read_metadata_with_dmtime(&self, file_path: impl AsRef<Path>) -> Result<(Vec<u8>, Option<OffsetDateTime>)> {
check_path_length(&file_path.as_ref().to_string_lossy().to_string())?;
@@ -624,15 +615,6 @@ pub async fn read_file_metadata(p: impl AsRef<Path>) -> Result<Metadata> {
Ok(meta)
}
pub async fn check_volume_exists(p: impl AsRef<Path>) -> Result<()> {
fs::metadata(&p).await.map_err(|e| match e.kind() {
ErrorKind::NotFound => Error::from(DiskError::VolumeNotFound),
ErrorKind::PermissionDenied => Error::from(DiskError::FileAccessDenied),
_ => Error::from(e),
})?;
Ok(())
}
fn skip_access_checks(p: impl AsRef<str>) -> bool {
let vols = [
super::RUSTFS_META_TMP_DELETED_BUCKET,
@@ -766,29 +748,22 @@ impl DiskAPI for LocalDisk {
}
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> {
let vol_path = self.get_bucket_path(volume)?;
let volume_dir = self.get_bucket_path(volume)?;
if !skip_access_checks(volume) {
check_volume_exists(&vol_path).await?;
if let Err(e) = utils::fs::access(&volume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
let fpath = self.get_object_path(volume, path)?;
let file_path = volume_dir.join(Path::new(&path));
check_path_length(file_path.to_string_lossy().to_string().as_str())?;
self.delete_file(&vol_path, &fpath, opt.recursive, opt.immediate).await?;
// if opt.recursive {
// let trash_path = self.get_object_path(RUSTFS_META_TMP_DELETED_BUCKET, Uuid::new_v4().to_string().as_str())?;
// fs::create_dir_all(&trash_path).await?;
// fs::rename(&fpath, &trash_path).await?;
// // TODO: immediate
// return Ok(());
// }
// fs::remove_file(fpath).await?;
self.delete_file(&volume_dir, &file_path, opt.recursive, opt.immediate)
.await?;
Ok(())
}
async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec<u8>) -> Result<()> {
let src_volume_dir = self.get_bucket_path(src_volume)?;
let dst_volume_dir = self.get_bucket_path(dst_volume)?;
@@ -871,156 +846,220 @@ impl DiskAPI for LocalDisk {
Ok(())
}
async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> {
let src_volume_path = self.get_bucket_path(src_volume)?;
let dst_volume_path = self.get_bucket_path(dst_volume)?;
let src_volume_dir = self.get_bucket_path(src_volume)?;
let dst_volume_dir = self.get_bucket_path(dst_volume)?;
if !skip_access_checks(src_volume) {
utils::fs::access(&src_volume_path).await.map_err(map_err_not_exists)?;
if let Err(e) = utils::fs::access(&src_volume_dir).await {
if os_is_not_exist(&e) {
return Err(Error::from(DiskError::VolumeNotFound));
} else if is_sys_err_io(&e) {
return Err(Error::from(DiskError::FaultyDisk));
}
return Err(Error::new(e));
}
}
if !skip_access_checks(dst_volume) {
utils::fs::access(&dst_volume_path).await.map_err(map_err_not_exists)?;
if let Err(e) = utils::fs::access(&dst_volume_dir).await {
if os_is_not_exist(&e) {
return Err(Error::from(DiskError::VolumeNotFound));
} else if is_sys_err_io(&e) {
return Err(Error::from(DiskError::FaultyDisk));
}
return Err(Error::new(e));
}
}
let srcp = self.get_object_path(src_volume, src_path)?;
let dstp = self.get_object_path(dst_volume, dst_path)?;
let src_is_dir = srcp.is_dir();
let dst_is_dir = dstp.is_dir();
if !src_is_dir && dst_is_dir || src_is_dir && !dst_is_dir {
let src_is_dir = has_suffix(&src_path, SLASH_SEPARATOR);
let dst_is_dir = has_suffix(&dst_path, SLASH_SEPARATOR);
if !(src_is_dir && dst_is_dir || !src_is_dir && !dst_is_dir) {
return Err(Error::from(DiskError::FileAccessDenied));
}
// TODO: check path length
let src_file_path = src_volume_dir.join(Path::new(&src_path));
check_path_length(src_file_path.to_string_lossy().to_string().as_str())?;
let dst_file_path = dst_volume_dir.join(Path::new(&dst_path));
check_path_length(dst_file_path.to_string_lossy().to_string().as_str())?;
if src_is_dir {
// TODO: remove dst_dir
}
let meta_op = match lstat(&src_file_path).await {
Ok(meta) => Some(meta),
Err(e) => {
if is_sys_err_io(&e) {
return Err(Error::new(DiskError::FaultyDisk));
}
fs::create_dir_all(dstp.parent().unwrap_or_else(|| Path::new("."))).await?;
let mut idx = 0;
loop {
if let Err(e) = fs::rename(&srcp, &dstp).await {
if e.kind() == ErrorKind::NotFound && idx == 0 {
idx += 1;
continue;
if !os_is_not_exist(&e) {
return Err(Error::new(e));
}
None
}
};
break;
if let Some(meta) = meta_op {
if !meta.is_dir() {
return Err(Error::new(DiskError::FileAccessDenied));
}
}
if let Err(e) = utils::fs::remove(&dst_file_path).await {
if is_sys_err_not_empty(&e) || is_sys_err_not_dir(&e) {
return Err(Error::new(DiskError::FileAccessDenied));
} else if is_sys_err_io(&e) {
return Err(Error::new(DiskError::FaultyDisk));
}
return Err(Error::new(e));
}
}
if let Some(dir_path) = srcp.parent() {
self.delete_file(&src_volume_path, &PathBuf::from(dir_path), false, false)
.await?;
if let Err(err) = os::rename_all(&src_file_path, &dst_file_path, &dst_volume_dir).await {
if let Some(e) = err.to_io_err() {
if is_sys_err_not_empty(&e) || is_sys_err_not_dir(&e) {
return Err(Error::new(DiskError::FileAccessDenied));
}
return Err(os_err_to_file_err(e));
}
return Err(err);
}
if let Some(parent) = src_file_path.parent() {
let _ = self.delete_file(&src_volume_dir, &parent.to_path_buf(), false, false).await;
}
Ok(())
}
async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, _file_size: usize) -> Result<FileWriter> {
let volpath = self.get_bucket_path(&volume)?;
// check exists
fs::metadata(&volpath).await.map_err(|e| match e.kind() {
ErrorKind::NotFound => Error::new(DiskError::VolumeNotFound),
_ => Error::new(e),
})?;
let fpath = self.get_object_path(volume, path)?;
debug!("CreateFile fpath: {:?}", fpath);
if let Some(_dir_path) = fpath.parent() {
fs::create_dir_all(&_dir_path).await?;
// TODO: use io.reader
async fn create_file(&self, origvolume: &str, volume: &str, path: &str, _file_size: usize) -> Result<FileWriter> {
if !origvolume.is_empty() {
let origvolume_dir = self.get_bucket_path(&origvolume)?;
if !skip_access_checks(origvolume) {
if let Err(e) = utils::fs::access(origvolume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
}
let file = File::create(&fpath).await?;
let volume_dir = self.get_bucket_path(&volume)?;
let file_path = volume_dir.join(Path::new(&path));
check_path_length(file_path.to_string_lossy().to_string().as_str())?;
Ok(FileWriter::Local(LocalFileWriter::new(file)))
// Ok(FileWriter::new(file))
// TODO: writeAllDirect io.copy
if let Some(parent) = file_path.parent() {
os::make_dir_all(parent, &volume_dir).await?;
}
let f = utils::fs::open_file(&file_path, O_CREATE | O_WRONLY)
.await
.map_err(os_err_to_file_err)?;
// let mut writer = BufWriter::new(file);
// io::copy(&mut r, &mut writer).await?;
Ok(FileWriter::Local(LocalFileWriter::new(f)))
// Ok(())
}
// async fn append_file(&self, volume: &str, path: &str, mut r: DuplexStream) -> Result<File> {
async fn append_file(&self, volume: &str, path: &str) -> Result<FileWriter> {
let p = self.get_object_path(volume, path)?;
if let Some(dir_path) = p.parent() {
fs::create_dir_all(&dir_path).await?;
}
// debug!("append_file open {} {:?}", self.id(), &p);
let file = File::options()
.read(true)
.create(true)
.write(true)
.append(true)
.open(&p)
.await?;
Ok(FileWriter::Local(LocalFileWriter::new(file)))
// Ok(FileWriter::new(file))
// let mut writer = BufWriter::new(file);
// io::copy(&mut r, &mut writer).await?;
// debug!("append_file end {} {}", self.id(), path);
// io::copy(&mut r, &mut file).await?;
// Ok(())
}
async fn read_file(&self, volume: &str, path: &str) -> Result<FileReader> {
let p = self.get_object_path(volume, path)?;
// debug!("read_file {:?}", &p);
let file = File::options().read(true).open(&p).await?;
Ok(FileReader::Local(LocalFileReader::new(file)))
// file.seek(SeekFrom::Start(offset as u64)).await?;
// let mut buffer = vec![0; length];
// let bytes_read = file.read(&mut buffer).await?;
// buffer.truncate(bytes_read);
// Ok((buffer, bytes_read))
}
async fn list_dir(&self, _origvolume: &str, volume: &str, _dir_path: &str, _count: i32) -> Result<Vec<String>> {
let p = self.get_bucket_path(volume)?;
let mut entries = fs::read_dir(&p).await?;
let mut volumes = Vec::new();
while let Some(entry) = entries.next_entry().await? {
if let Ok(_metadata) = entry.metadata().await {
// if !metadata.is_dir() {
// continue;
// }
let name = entry.file_name().to_string_lossy().to_string();
// let created = match metadata.created() {
// Ok(md) => OffsetDateTime::from(md),
// Err(_) => return Err(Error::msg("Not supported created on this platform")),
// };
volumes.push(name);
let volume_dir = self.get_bucket_path(&volume)?;
if !skip_access_checks(volume) {
if let Err(e) = utils::fs::access(&volume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
Ok(volumes)
let file_path = volume_dir.join(Path::new(&path));
check_path_length(file_path.to_string_lossy().to_string().as_str())?;
let f = self.open_file(file_path, O_CREATE | O_APPEND | O_WRONLY, volume_dir).await?;
Ok(FileWriter::Local(LocalFileWriter::new(f)))
}
// TODO: io verifier
async fn read_file(&self, volume: &str, path: &str) -> Result<FileReader> {
let volume_dir = self.get_bucket_path(&volume)?;
if !skip_access_checks(volume) {
if let Err(e) = utils::fs::access(&volume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
let file_path = volume_dir.join(Path::new(&path));
check_path_length(file_path.to_string_lossy().to_string().as_str())?;
let f = self.open_file(file_path, O_RDONLY, volume_dir).await.map_err(|err| {
if let Some(e) = err.to_io_err() {
if os_is_not_exist(&e) {
Error::new(DiskError::FileNotFound)
} else if os_is_permission(&e) {
Error::new(DiskError::FileAccessDenied)
} else if is_sys_err_not_dir(&e) {
Error::new(DiskError::FileAccessDenied)
} else if is_sys_err_io(&e) {
Error::new(DiskError::FaultyDisk)
} else if is_sys_err_too_many_files(&e) {
Error::new(DiskError::TooManyOpenFiles)
} else {
Error::new(e)
}
} else {
err
}
})?;
Ok(FileReader::Local(LocalFileReader::new(f)))
}
async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> Result<Vec<String>> {
if !origvolume.is_empty() {
let origvolume_dir = self.get_bucket_path(&origvolume)?;
if !skip_access_checks(origvolume) {
if let Err(e) = utils::fs::access(origvolume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
}
let volume_dir = self.get_bucket_path(&volume)?;
let dir_path_abs = volume_dir.join(Path::new(&dir_path));
let entries = match os::read_dir(&dir_path_abs, count).await {
Ok(res) => res,
Err(e) => {
if DiskError::FileNotFound.is(&e) && !skip_access_checks(&volume) {
if let Err(e) = utils::fs::access(&volume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
return Err(e);
}
};
Ok(entries)
}
// TODO: io.writer
async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>> {
let mut entries = self.list_dir("", &opts.bucket, &opts.base_dir, -1).await?;
let mut entries = match self.list_dir("", &opts.bucket, &opts.base_dir, -1).await {
Ok(res) => res,
Err(e) => {
if !DiskError::VolumeNotFound.is(&e) && !DiskError::FileNotFound.is(&e) {
error!("list_dir err {:?}", &e);
}
if opts.report_notfound && DiskError::FileNotFound.is(&e) {
return Err(Error::new(DiskError::FileNotFound));
}
return Ok(Vec::new());
}
};
if entries.is_empty() {
return Ok(Vec::new());
}
entries.sort();
@@ -1051,12 +1090,7 @@ impl DiskAPI for LocalDisk {
let fpath = self.get_object_path(bucket, format!("{}/{}", &meta.name, STORAGE_FORMAT_FILE).as_str())?;
let (fdata, _) = match self.read_metadata_with_dmtime(&fpath).await {
Ok(res) => res,
Err(_) => (Vec::new(), None),
};
meta.metadata = fdata;
meta.metadata = self.read_metadata(&fpath).await.unwrap_or_default();
metas.push(meta);
}
@@ -1073,107 +1107,196 @@ impl DiskAPI for LocalDisk {
dst_volume: &str,
dst_path: &str,
) -> Result<RenameDataResp> {
let src_volume_path = self.get_bucket_path(src_volume)?;
let src_volume_dir = self.get_bucket_path(src_volume)?;
if !skip_access_checks(src_volume) {
check_volume_exists(&src_volume_path).await?;
if let Err(e) = utils::fs::access(&src_volume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
let dst_volume_path = self.get_bucket_path(dst_volume)?;
let dst_volume_dir = self.get_bucket_path(dst_volume)?;
if !skip_access_checks(dst_volume) {
check_volume_exists(&dst_volume_path).await?;
if let Err(e) = utils::fs::access(&dst_volume_dir).await {
return Err(convert_access_error(e, DiskError::VolumeAccessDenied));
}
}
// xl.meta路径
let src_file_path = self.get_object_path(src_volume, format!("{}/{}", &src_path, super::STORAGE_FORMAT_FILE).as_str())?;
let dst_file_path = self.get_object_path(dst_volume, format!("{}/{}", &dst_path, super::STORAGE_FORMAT_FILE).as_str())?;
let src_file_path = src_volume_dir.join(Path::new(format!("{}/{}", &src_path, super::STORAGE_FORMAT_FILE).as_str()));
let dst_file_path = dst_volume_dir.join(Path::new(format!("{}/{}", &dst_path, super::STORAGE_FORMAT_FILE).as_str()));
// data_dir 路径
let (src_data_path, dst_data_path) = {
let mut data_dir = String::new();
if !fi.is_remote() {
data_dir = utils::path::retain_slash(fi.data_dir.unwrap_or(Uuid::nil()).to_string().as_str());
}
let has_data_dir_path = {
let has_data_dir = {
if !fi.is_remote() {
if let Some(dir) = fi.data_dir {
Some(utils::path::retain_slash(dir.to_string().as_str()))
} else {
None
}
} else {
None
}
};
if !data_dir.is_empty() {
let src_data_path = self.get_object_path(
src_volume,
if let Some(data_dir) = has_data_dir {
let src_data_path = src_volume_dir.join(Path::new(
utils::path::retain_slash(format!("{}/{}", &src_path, data_dir).as_str()).as_str(),
)?;
let dst_data_path = self.get_object_path(dst_volume, format!("{}/{}", &dst_path, data_dir).as_str())?;
));
let dst_data_path = dst_volume_dir.join(Path::new(
utils::path::retain_slash(format!("{}/{}", &dst_path, data_dir).as_str()).as_str(),
));
(src_data_path, dst_data_path)
Some((src_data_path, dst_data_path))
} else {
(PathBuf::new(), PathBuf::new())
None
}
};
// 读旧xl.meta
let mut meta = FileMeta::new();
check_path_length(&src_file_path.to_string_lossy().to_string().as_str())?;
check_path_length(&dst_file_path.to_string_lossy().to_string().as_str())?;
let (dst_buf, _) = read_file_exists(&dst_file_path).await?;
if !dst_buf.is_empty() {
// 有旧文件,加载
match meta.unmarshal_msg(&dst_buf) {
Ok(_) => {}
Err(_) => meta = FileMeta::new(),
// 读旧xl.meta
let has_dst_buf = match utils::fs::read_file(&dst_file_path).await {
Ok(res) => Some(res),
Err(e) => {
if is_sys_err_not_dir(&e) && !cfg!(target_os = "windows") {
return Err(Error::new(DiskError::FileAccessDenied));
}
if !os_is_not_exist(&e) {
return Err(os_err_to_file_err(e));
}
None
}
};
// let current_data_path = dst_volume_dir.join(Path::new(&dst_path));
let mut xlmeta = FileMeta::new();
if let Some(dst_buf) = has_dst_buf.as_ref() {
if FileMeta::is_xl_format(&dst_buf) {
if let Ok(nmeta) = FileMeta::load(&dst_buf) {
xlmeta = nmeta
}
}
}
let mut skip_parent = dst_volume_path.clone();
if !dst_buf.is_empty() {
skip_parent = PathBuf::from(&dst_file_path.parent().unwrap_or(Path::new("/")));
let mut skip_parent = dst_volume_dir.clone();
if has_dst_buf.as_ref().is_some() {
if let Some(parent) = dst_file_path.parent() {
skip_parent = parent.to_path_buf();
}
}
// 查找版本是否已存在
let old_data_dir = meta
.find_version(fi.version_id)
.map(|(_, version)| {
version
.get_data_dir()
.filter(|data_dir| meta.shard_data_dir_count(&fi.version_id, &Some(data_dir.clone())) == 0)
})
.unwrap_or_default();
// TODO: Healing
// 添加版本,写入xl.meta文件
meta.add_version(fi.clone())?;
let has_old_data_dir = {
if let Ok((_, ver)) = xlmeta.find_version(fi.version_id) {
let has_data_dir = ver.get_data_dir();
if let Some(data_dir) = has_data_dir {
if xlmeta.shard_data_dir_count(&fi.version_id, &Some(data_dir.clone())) == 0 {
// TODO: Healing
// remove inlinedata\
Some(data_dir)
} else {
None
}
} else {
None
}
} else {
None
}
};
let fm_data = meta.marshal_msg()?;
xlmeta.add_version(fi.clone())?;
// 写入xl.meta
self.write_all(src_volume, format!("{}/{}", &src_path, super::STORAGE_FORMAT_FILE).as_str(), fm_data)
.await?;
let no_inline = src_data_path.has_root() && fi.data.is_none() && fi.size > 0;
if no_inline {
self.rename_all(&src_data_path, &dst_data_path, &skip_parent).await?;
if xlmeta.versions.len() <= 10 {
// TODO: Sign
}
// warn!("old_data_dir {:?}", old_data_dir);
// 有旧目录,把old xl.meta存到旧目录里
if old_data_dir.is_some() {
self.write_all(
&dst_volume,
format!("{}/{}/{}", &dst_path, &old_data_dir.unwrap().to_string(), super::STORAGE_FORMAT_FILE).as_str(),
dst_buf,
)
.await?;
let new_dst_buf = xlmeta.marshal_msg()?;
self.write_all(&src_volume, format!("{}/{}", &src_path, super::STORAGE_FORMAT_FILE).as_str(), new_dst_buf)
.await
.map_err(|err| {
if let Some(e) = err.to_io_err() {
os_err_to_file_err(e)
} else {
err
}
})?;
if let Some((src_data_path, dst_data_path)) = has_data_dir_path.as_ref() {
let no_inline = fi.data.is_none() && fi.size > 0;
if no_inline {
if let Err(err) = os::rename_all(&src_data_path, &dst_data_path, &skip_parent).await {
let _ = self.delete_file(&dst_volume_dir, &dst_data_path, false, false).await;
return Err({
if let Some(e) = err.to_io_err() {
os_err_to_file_err(e)
} else {
err
}
});
}
}
}
if let Err(e) = self.rename_all(&src_file_path, &dst_file_path, &skip_parent).await {
// 如果 失败删除目标目录
let _ = self.delete_file(&dst_volume_path, &dst_data_path, false, false).await;
return Err(e);
if let Some(old_data_dir) = has_old_data_dir {
// preserve current xl.meta inside the oldDataDir.
if let Some(dst_buf) = has_dst_buf {
if let Err(err) = self
.write_all_private(
&dst_volume,
format!("{}/{}/{}", &dst_path, &old_data_dir.to_string(), super::STORAGE_FORMAT_FILE).as_str(),
&dst_buf,
true,
&skip_parent,
)
.await
{
return Err({
if let Some(e) = err.to_io_err() {
os_err_to_file_err(e)
} else {
err
}
});
}
}
}
if src_volume != super::RUSTFS_META_MULTIPART_BUCKET {
fs::remove_dir(&src_file_path.parent().unwrap()).await?;
} else {
self.delete_file(&src_volume_path, &PathBuf::from(src_file_path.parent().unwrap()), true, false)
.await?;
if let Err(err) = os::rename_all(&src_file_path, &dst_file_path, &skip_parent).await {
if let Some((_, dst_data_path)) = has_data_dir_path.as_ref() {
let _ = self.delete_file(&dst_volume_dir, &dst_data_path, false, false).await;
}
return Err({
if let Some(e) = err.to_io_err() {
os_err_to_file_err(e)
} else {
err
}
});
}
if let Some(src_file_path_parent) = src_file_path.parent() {
if src_volume != super::RUSTFS_META_MULTIPART_BUCKET {
let _ = utils::fs::remove(src_file_path_parent).await;
} else {
let _ = self
.delete_file(&dst_volume_dir, &src_file_path_parent.to_path_buf(), true, false)
.await;
}
}
Ok(RenameDataResp {
old_data_dir: old_data_dir,
old_data_dir: has_old_data_dir,
sign: None, // TODO:
})
}
@@ -1194,12 +1317,19 @@ impl DiskAPI for LocalDisk {
return Err(Error::msg("Invalid arguments specified"));
}
let p = self.get_bucket_path(volume)?;
let volume_dir = self.get_bucket_path(volume)?;
if let Err(e) = utils::fs::access(&p).await {
if let Err(e) = utils::fs::access(&volume_dir).await {
if os_is_not_exist(&e) {
os::make_dir_all(&p, self.root.as_path()).await?;
os::make_dir_all(&volume_dir, self.root.as_path()).await?;
}
if os_is_permission(&e) {
return Err(Error::new(DiskError::DiskAccessDenied));
} else if is_sys_err_io(&e) {
return Err(Error::new(DiskError::FaultyDisk));
}
return Err(Error::new(e));
}
Err(Error::from(DiskError::VolumeExists))
@@ -1207,7 +1337,7 @@ impl DiskAPI for LocalDisk {
async fn list_volumes(&self) -> Result<Vec<VolumeInfo>> {
let mut volumes = Vec::new();
let entries = os::read_dir(&self.root, 0).await.map_err(|e| {
let entries = os::read_dir(&self.root, -1).await.map_err(|e| {
if DiskError::FileAccessDenied.is(&e) {
Error::new(DiskError::DiskAccessDenied)
} else if DiskError::FileNotFound.is(&e) {
@@ -1231,17 +1361,27 @@ impl DiskAPI for LocalDisk {
Ok(volumes)
}
async fn stat_volume(&self, volume: &str) -> Result<VolumeInfo> {
let p = self.get_bucket_path(volume)?;
let m = read_file_metadata(&p).await?;
let modtime = match m.modified() {
Ok(md) => Some(OffsetDateTime::from(md)),
Err(_) => {
warn!("Not supported modified on this platform");
None
let volume_dir = self.get_bucket_path(volume)?;
let meta = match utils::fs::lstat(&volume_dir).await {
Ok(res) => res,
Err(e) => {
if os_is_not_exist(&e) {
return Err(Error::new(DiskError::VolumeNotFound));
} else if os_is_permission(&e) {
return Err(Error::new(DiskError::DiskAccessDenied));
} else if is_sys_err_io(&e) {
return Err(Error::new(DiskError::FaultyDisk));
} else {
return Err(Error::new(e));
}
}
};
let modtime = match meta.modified() {
Ok(md) => Some(OffsetDateTime::from(md)),
Err(_) => None,
};
Ok(VolumeInfo {
name: volume.to_string(),
created: modtime,
+2 -8
View File
@@ -68,18 +68,12 @@ pub async fn make_dir_all(path: impl AsRef<Path>, base_dir: impl AsRef<Path>) ->
}
// read_dir count read limit. when count == 0 unlimit.
pub async fn read_dir(path: impl AsRef<Path>, count: usize) -> Result<Vec<String>> {
pub async fn read_dir(path: impl AsRef<Path>, count: i32) -> Result<Vec<String>> {
let mut entries = fs::read_dir(path.as_ref()).await?;
let mut volumes = Vec::new();
let mut count: i32 = {
if count == 0 {
-1
} else {
count as i32
}
};
let mut count = count;
while let Some(entry) = entries.next_entry().await? {
let name = entry.file_name().to_string_lossy().to_string();
+4
View File
@@ -130,3 +130,7 @@ pub async fn mkdir(path: impl AsRef<Path>) -> io::Result<()> {
pub async fn rename(from: impl AsRef<Path>, to: impl AsRef<Path>) -> io::Result<()> {
fs::rename(from, to).await
}
pub async fn read_file(path: impl AsRef<Path>) -> io::Result<Vec<u8>> {
fs::read(path.as_ref()).await
}