Merge pull request #63 from rustfs/quorum

Quorum
This commit is contained in:
weisd
2024-09-26 17:26:13 +08:00
committed by GitHub
28 changed files with 2754 additions and 1436 deletions
+1
View File
@@ -2,3 +2,4 @@
.DS_Store .DS_Store
.idea .idea
.vscode .vscode
/test
@@ -1,10 +1,9 @@
// automatically generated by the FlatBuffers compiler, do not modify // automatically generated by the FlatBuffers compiler, do not modify
// @generated // @generated
use core::mem;
use core::cmp::Ordering; use core::cmp::Ordering;
use core::mem;
extern crate flatbuffers; extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow}; use self::flatbuffers::{EndianScalar, Follow};
@@ -12,112 +11,114 @@ use self::flatbuffers::{EndianScalar, Follow};
#[allow(unused_imports, dead_code)] #[allow(unused_imports, dead_code)]
pub mod models { pub mod models {
use core::mem; use core::cmp::Ordering;
use core::cmp::Ordering; use core::mem;
extern crate flatbuffers; extern crate flatbuffers;
use self::flatbuffers::{EndianScalar, Follow}; use self::flatbuffers::{EndianScalar, Follow};
pub enum PingBodyOffset {} pub enum PingBodyOffset {}
#[derive(Copy, Clone, PartialEq)] #[derive(Copy, Clone, PartialEq)]
pub struct PingBody<'a> { pub struct PingBody<'a> {
pub _tab: flatbuffers::Table<'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,
} }
}
}
pub struct PingBodyBuilder<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> { impl<'a> flatbuffers::Follow<'a> for PingBody<'a> {
fbb_: &'b mut flatbuffers::FlatBufferBuilder<'a, A>, type Inner = PingBody<'a>;
start_: flatbuffers::WIPOffset<flatbuffers::TableUnfinishedWIPOffset>, #[inline]
} unsafe fn follow(buf: &'a [u8], loc: usize) -> Self::Inner {
impl<'a: 'b, 'b, A: flatbuffers::Allocator + 'a> PingBodyBuilder<'a, 'b, A> { Self {
#[inline] _tab: flatbuffers::Table::new(buf, loc),
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<'_> { impl<'a> PingBody<'a> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { pub const VT_PAYLOAD: flatbuffers::VOffsetT = 4;
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
+57
View File
@@ -88,6 +88,20 @@ message DeleteResponse {
optional string error_info = 2; optional string error_info = 2;
} }
message RenamePartRequst {
string disk = 1;
string src_volume = 2;
string src_path = 3;
string dst_volume = 4;
string dst_path = 5;
bytes meta = 6;
}
message RenamePartResponse {
bool success = 1;
optional string error_info = 2;
}
message RenameFileRequst { message RenameFileRequst {
string disk = 1; string disk = 1;
string src_volume = 2; string src_volume = 2;
@@ -219,6 +233,30 @@ message StatVolumeResponse {
optional string error_info = 3; optional string error_info = 3;
} }
message DeletePathsRequest {
string disk = 1;
string volume = 2;
repeated string paths = 3;
}
message DeletePathsResponse {
bool success = 1;
optional string error_info = 2;
}
message UpdateMetadataRequest {
string disk = 1;
string volume = 2;
string path = 3;
string file_info = 4;
string opts = 5;
}
message UpdateMetadataResponse {
bool success = 1;
optional string error_info = 2;
}
message WriteMetadataRequest { message WriteMetadataRequest {
string disk = 1; // indicate which one in the disks string disk = 1; // indicate which one in the disks
string volume = 2; string volume = 2;
@@ -258,6 +296,21 @@ message ReadXLResponse {
optional string error_info = 3; optional string error_info = 3;
} }
message DeleteVersionRequest {
string disk = 1;
string volume = 2;
string path = 3;
string file_info = 4;
bool force_del_marker = 5;
string opts = 6;
}
message DeleteVersionResponse {
bool success = 1;
string raw_file_info = 2;
optional string error_info = 3;
}
message DeleteVersionsRequest { message DeleteVersionsRequest {
string disk = 1; string disk = 1;
string volume = 2; string volume = 2;
@@ -307,6 +360,7 @@ service NodeService {
rpc ReadAll(ReadAllRequest) returns (ReadAllResponse) {}; rpc ReadAll(ReadAllRequest) returns (ReadAllResponse) {};
rpc WriteAll(WriteAllRequest) returns (WriteAllResponse) {}; rpc WriteAll(WriteAllRequest) returns (WriteAllResponse) {};
rpc Delete(DeleteRequest) returns (DeleteResponse) {}; rpc Delete(DeleteRequest) returns (DeleteResponse) {};
rpc RenamePart(RenamePartRequst) returns (RenamePartResponse) {};
rpc RenameFile(RenameFileRequst) returns (RenameFileResponse) {}; rpc RenameFile(RenameFileRequst) returns (RenameFileResponse) {};
rpc Write(WriteRequest) returns (WriteResponse) {}; rpc Write(WriteRequest) returns (WriteResponse) {};
// rpc Append(AppendRequest) returns (AppendResponse) {}; // rpc Append(AppendRequest) returns (AppendResponse) {};
@@ -318,9 +372,12 @@ service NodeService {
rpc MakeVolume(MakeVolumeRequest) returns (MakeVolumeResponse) {}; rpc MakeVolume(MakeVolumeRequest) returns (MakeVolumeResponse) {};
rpc ListVolumes(ListVolumesRequest) returns (ListVolumesResponse) {}; rpc ListVolumes(ListVolumesRequest) returns (ListVolumesResponse) {};
rpc StatVolume(StatVolumeRequest) returns (StatVolumeResponse) {}; rpc StatVolume(StatVolumeRequest) returns (StatVolumeResponse) {};
rpc DeletePaths(DeletePathsRequest) returns (DeletePathsResponse) {};
rpc UpdateMetadata(UpdateMetadataRequest) returns (UpdateMetadataResponse) {};
rpc WriteMetadata(WriteMetadataRequest) returns (WriteMetadataResponse) {}; rpc WriteMetadata(WriteMetadataRequest) returns (WriteMetadataResponse) {};
rpc ReadVersion(ReadVersionRequest) returns (ReadVersionResponse) {}; rpc ReadVersion(ReadVersionRequest) returns (ReadVersionResponse) {};
rpc ReadXL(ReadXLRequest) returns (ReadXLResponse) {}; rpc ReadXL(ReadXLRequest) returns (ReadXLResponse) {};
rpc DeleteVersion(DeleteVersionRequest) returns (DeleteVersionResponse) {};
rpc DeleteVersions(DeleteVersionsRequest) returns (DeleteVersionsResponse) {}; rpc DeleteVersions(DeleteVersionsRequest) returns (DeleteVersionsResponse) {};
rpc ReadMultiple(ReadMultipleRequest) returns (ReadMultipleResponse) {}; rpc ReadMultiple(ReadMultipleRequest) returns (ReadMultipleResponse) {};
rpc DeleteVolume(DeleteVolumeRequest) returns (DeleteVolumeResponse) {}; rpc DeleteVolume(DeleteVolumeRequest) returns (DeleteVolumeResponse) {};
+323 -16
View File
@@ -1,39 +1,105 @@
use crate::error::{Error, Result}; use std::io::{self, ErrorKind};
use crate::{
error::{Error, Result},
quorum::CheckErrorFn,
};
// DiskError == StorageErr
#[derive(Debug, thiserror::Error)] #[derive(Debug, thiserror::Error)]
pub enum DiskError { pub enum DiskError {
#[error("maximum versions exceeded, please delete few versions to proceed")]
MaxVersionsExceeded,
#[error("unexpected error")]
Unexpected,
#[error("corrupted format")]
CorruptedFormat,
#[error("corrupted backend")]
CorruptedBackend,
#[error("unformatted disk error")]
UnformattedDisk,
#[error("inconsistent drive found")]
InconsistentDisk,
#[error("drive does not support O_DIRECT")]
UnsupportedDisk,
#[error("drive path full")]
DiskFull,
#[error("disk not a dir")]
DiskNotDir,
#[error("disk not found")]
DiskNotFound,
#[error("drive still did not complete the request")]
DiskOngoingReq,
#[error("drive is part of root drive, will not be used")]
DriveIsRoot,
#[error("remote drive is faulty")]
FaultyRemoteDisk,
#[error("drive is faulty")]
FaultyDisk,
#[error("drive access denied")]
DiskAccessDenied,
#[error("file not found")] #[error("file not found")]
FileNotFound, FileNotFound,
#[error("file version not found")] #[error("file version not found")]
FileVersionNotFound, FileVersionNotFound,
#[error("disk not found")] #[error("too many open files, please increase 'ulimit -n'")]
DiskNotFound, TooManyOpenFiles,
#[error("disk access denied")] #[error("file name too long")]
FileAccessDenied, FileNameTooLong,
#[error("InconsistentDisk")]
InconsistentDisk,
#[error("volume already exists")] #[error("volume already exists")]
VolumeExists, VolumeExists,
#[error("unformatted disk error")] #[error("not of regular file type")]
UnformattedDisk, IsNotRegular,
#[error("unsupport disk")] #[error("path not found")]
UnsupportedDisk, PathNotFound,
#[error("disk not a dir")]
DiskNotDir,
#[error("volume not found")] #[error("volume not found")]
VolumeNotFound, VolumeNotFound,
#[error("volume not empty")] #[error("volume is not empty")]
VolumeNotEmpty, VolumeNotEmpty,
#[error("volume access denied")]
VolumeAccessDenied,
#[error("disk access denied")]
FileAccessDenied,
#[error("file is corrupted")]
FileCorrupt,
#[error("bit-rot hash algorithm is invalid")]
BitrotHashAlgoInvalid,
#[error("Rename across devices not allowed, please fix your backend configuration")]
CrossDeviceLink,
#[error("less data available than what was requested")]
LessData,
#[error("more data was sent than what was advertised")]
MoreData,
} }
impl DiskError { impl DiskError {
@@ -96,3 +162,244 @@ impl PartialEq for DiskError {
core::mem::discriminant(self) == core::mem::discriminant(other) core::mem::discriminant(self) == core::mem::discriminant(other)
} }
} }
impl CheckErrorFn for DiskError {
fn is(&self, e: &Error) -> bool {
self.is(e)
}
}
pub fn clone_disk_err(e: &DiskError) -> Error {
match e {
DiskError::MaxVersionsExceeded => Error::new(DiskError::MaxVersionsExceeded),
DiskError::Unexpected => Error::new(DiskError::Unexpected),
DiskError::CorruptedFormat => Error::new(DiskError::CorruptedFormat),
DiskError::CorruptedBackend => Error::new(DiskError::CorruptedBackend),
DiskError::UnformattedDisk => Error::new(DiskError::UnformattedDisk),
DiskError::InconsistentDisk => Error::new(DiskError::InconsistentDisk),
DiskError::UnsupportedDisk => Error::new(DiskError::UnsupportedDisk),
DiskError::DiskFull => Error::new(DiskError::DiskFull),
DiskError::DiskNotDir => Error::new(DiskError::DiskNotDir),
DiskError::DiskNotFound => Error::new(DiskError::DiskNotFound),
DiskError::DiskOngoingReq => Error::new(DiskError::DiskOngoingReq),
DiskError::DriveIsRoot => Error::new(DiskError::DriveIsRoot),
DiskError::FaultyRemoteDisk => Error::new(DiskError::FaultyRemoteDisk),
DiskError::FaultyDisk => Error::new(DiskError::FaultyDisk),
DiskError::DiskAccessDenied => Error::new(DiskError::DiskAccessDenied),
DiskError::FileNotFound => Error::new(DiskError::FileNotFound),
DiskError::FileVersionNotFound => Error::new(DiskError::FileVersionNotFound),
DiskError::TooManyOpenFiles => Error::new(DiskError::TooManyOpenFiles),
DiskError::FileNameTooLong => Error::new(DiskError::FileNameTooLong),
DiskError::VolumeExists => Error::new(DiskError::VolumeExists),
DiskError::IsNotRegular => Error::new(DiskError::IsNotRegular),
DiskError::PathNotFound => Error::new(DiskError::PathNotFound),
DiskError::VolumeNotFound => Error::new(DiskError::VolumeNotFound),
DiskError::VolumeNotEmpty => Error::new(DiskError::VolumeNotEmpty),
DiskError::VolumeAccessDenied => Error::new(DiskError::VolumeAccessDenied),
DiskError::FileAccessDenied => Error::new(DiskError::FileAccessDenied),
DiskError::FileCorrupt => Error::new(DiskError::FileCorrupt),
DiskError::BitrotHashAlgoInvalid => Error::new(DiskError::BitrotHashAlgoInvalid),
DiskError::CrossDeviceLink => Error::new(DiskError::CrossDeviceLink),
DiskError::LessData => Error::new(DiskError::LessData),
DiskError::MoreData => Error::new(DiskError::MoreData),
}
}
pub fn os_err_to_file_err(e: io::Error) -> Error {
match e.kind() {
io::ErrorKind::NotFound => Error::new(DiskError::FileNotFound),
io::ErrorKind::PermissionDenied => Error::new(DiskError::FileAccessDenied),
// io::ErrorKind::ConnectionRefused => todo!(),
// io::ErrorKind::ConnectionReset => todo!(),
// io::ErrorKind::HostUnreachable => todo!(),
// io::ErrorKind::NetworkUnreachable => todo!(),
// io::ErrorKind::ConnectionAborted => todo!(),
// io::ErrorKind::NotConnected => todo!(),
// io::ErrorKind::AddrInUse => todo!(),
// io::ErrorKind::AddrNotAvailable => todo!(),
// io::ErrorKind::NetworkDown => todo!(),
// io::ErrorKind::BrokenPipe => todo!(),
// io::ErrorKind::AlreadyExists => todo!(),
// io::ErrorKind::WouldBlock => todo!(),
// io::ErrorKind::NotADirectory => DiskError::FileNotFound,
// io::ErrorKind::IsADirectory => DiskError::FileNotFound,
// io::ErrorKind::DirectoryNotEmpty => DiskError::VolumeNotEmpty,
// io::ErrorKind::ReadOnlyFilesystem => todo!(),
// io::ErrorKind::FilesystemLoop => todo!(),
// io::ErrorKind::StaleNetworkFileHandle => todo!(),
// io::ErrorKind::InvalidInput => todo!(),
// io::ErrorKind::InvalidData => todo!(),
// io::ErrorKind::TimedOut => todo!(),
// io::ErrorKind::WriteZero => todo!(),
// io::ErrorKind::StorageFull => DiskError::DiskFull,
// io::ErrorKind::NotSeekable => todo!(),
// io::ErrorKind::FilesystemQuotaExceeded => todo!(),
// io::ErrorKind::FileTooLarge => todo!(),
// io::ErrorKind::ResourceBusy => todo!(),
// io::ErrorKind::ExecutableFileBusy => todo!(),
// io::ErrorKind::Deadlock => todo!(),
// io::ErrorKind::CrossesDevices => todo!(),
// io::ErrorKind::TooManyLinks =>DiskError::TooManyOpenFiles,
// io::ErrorKind::InvalidFilename => todo!(),
// io::ErrorKind::ArgumentListTooLong => todo!(),
// io::ErrorKind::Interrupted => todo!(),
// io::ErrorKind::Unsupported => todo!(),
// io::ErrorKind::UnexpectedEof => todo!(),
// io::ErrorKind::OutOfMemory => todo!(),
// io::ErrorKind::Other => todo!(),
// TODO: 把不支持的king用字符串处理
_ => Error::new(e),
}
}
pub fn is_sys_err_no_space(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 28;
}
false
}
pub fn is_sys_err_invalid_arg(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 22;
}
false
}
pub fn is_sys_err_io(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 5;
}
false
}
pub fn is_sys_err_is_dir(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 21;
}
false
}
pub fn is_sys_err_not_dir(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 20;
}
false
}
pub fn is_sys_err_too_long(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 63;
}
false
}
pub fn is_sys_err_too_many_symlinks(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 62;
}
false
}
pub fn is_sys_err_not_empty(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
if no == 66 {
return true;
}
if cfg!(target_os = "solaris") && no == 17 {
return true;
}
if cfg!(target_os = "windows") && no == 145 {
return true;
}
}
false
}
pub fn is_sys_err_path_not_found(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
if cfg!(target_os = "windows") {
if no == 3 {
return true;
}
} else {
if no == 2 {
return true;
}
}
}
false
}
pub fn is_sys_err_handle_invalid(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
if cfg!(target_os = "windows") {
if no == 6 {
return true;
}
} else {
return false;
}
}
false
}
pub fn is_sys_err_cross_device(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 18;
}
false
}
pub fn is_sys_err_too_many_files(e: &io::Error) -> bool {
if let Some(no) = e.raw_os_error() {
return no == 23 || no == 24;
}
false
}
pub fn os_is_not_exist(e: &io::Error) -> bool {
e.kind() == ErrorKind::NotFound
}
pub fn os_is_permission(e: &io::Error) -> bool {
if e.kind() == ErrorKind::PermissionDenied {
return true;
}
if let Some(no) = e.raw_os_error() {
if no == 30 {
return true;
}
}
false
}
pub fn os_is_exist(e: &io::Error) -> bool {
e.kind() == ErrorKind::AlreadyExists
}
// map_err_not_exists
pub fn map_err_not_exists(e: io::Error) -> Error {
if os_is_not_exist(&e) {
return Error::new(DiskError::VolumeNotEmpty);
} else if is_sys_err_io(&e) {
return Error::new(DiskError::FaultyDisk);
}
Error::new(e)
}
pub fn convert_access_error(e: io::Error, per_err: DiskError) -> Error {
if os_is_not_exist(&e) {
return Error::new(DiskError::VolumeNotEmpty);
} else if is_sys_err_io(&e) {
return Error::new(DiskError::FaultyDisk);
} else if os_is_permission(&e) {
return Error::new(per_err);
}
Error::new(e)
}
+49
View File
@@ -182,6 +182,55 @@ impl FormatV3 {
Err(Error::msg(format!("disk id not found {}", disk_id))) Err(Error::msg(format!("disk id not found {}", disk_id)))
} }
pub fn check_other(&self, other: &FormatV3) -> Result<()> {
let mut tmp = other.clone();
let this = tmp.erasure.this;
tmp.erasure.this = Uuid::nil();
if self.erasure.sets.len() != other.erasure.sets.len() {
return Err(Error::from_string(format!(
"Expected number of sets {}, got {}",
self.erasure.sets.len(),
other.erasure.sets.len()
)));
}
for i in 0..self.erasure.sets.len() {
if self.erasure.sets[i].len() != other.erasure.sets[i].len() {
return Err(Error::from_string(format!(
"Each set should be of same size, expected {}, got {}",
self.erasure.sets[i].len(),
other.erasure.sets[i].len()
)));
}
for j in 0..self.erasure.sets[i].len() {
if self.erasure.sets[i][j] != self.erasure.sets[i][j] {
return Err(Error::from_string(format!(
"UUID on positions {}:{} do not match with, expected {:?} got {:?}: (%w)",
i,
j,
self.erasure.sets[i][j].to_string(),
other.erasure.sets[i][j].to_string(),
)));
}
}
}
for i in 0..tmp.erasure.sets.len() {
for j in 0..tmp.erasure.sets[i].len() {
if this == tmp.erasure.sets[i][j] {
return Ok(());
}
}
}
Err(Error::msg(format!(
"DriveID {:?} not found in any drive sets {:?}",
this, other.erasure.sets
)))
}
} }
#[cfg(test)] #[cfg(test)]
+568 -123
View File
@@ -1,16 +1,23 @@
use super::error::{is_sys_err_io, is_sys_err_not_empty, is_sys_err_too_many_files, os_is_not_exist, os_is_permission};
use super::{endpoint::Endpoint, error::DiskError, format::FormatV3}; use super::{endpoint::Endpoint, error::DiskError, format::FormatV3};
use super::{ use super::{
DeleteOptions, DiskAPI, FileInfoVersions, FileReader, FileWriter, MetaCacheEntry, ReadMultipleReq, ReadMultipleResp, os, DeleteOptions, DiskAPI, DiskLocation, FileInfoVersions, FileReader, FileWriter, MetaCacheEntry, ReadMultipleReq,
ReadOptions, RenameDataResp, VolumeInfo, WalkDirOptions, ReadMultipleResp, ReadOptions, RenameDataResp, UpdateMetadataOpts, VolumeInfo, WalkDirOptions,
}; };
use crate::disk::error::{
convert_access_error, is_sys_err_handle_invalid, is_sys_err_invalid_arg, is_sys_err_is_dir, is_sys_err_not_dir,
map_err_not_exists, os_err_to_file_err,
};
use crate::disk::os::check_path_length;
use crate::disk::{LocalFileReader, LocalFileWriter, STORAGE_FORMAT_FILE}; use crate::disk::{LocalFileReader, LocalFileWriter, STORAGE_FORMAT_FILE};
use crate::utils::fs::{lstat, O_RDONLY};
use crate::utils::path::{has_suffix, SLASH_SEPARATOR};
use crate::{ use crate::{
error::{Error, Result}, error::{Error, Result},
file_meta::FileMeta, file_meta::FileMeta,
store_api::{FileInfo, RawFileInfo}, store_api::{FileInfo, RawFileInfo},
utils, utils,
}; };
use bytes::Bytes;
use path_absolutize::Absolutize; use path_absolutize::Absolutize;
use std::{ use std::{
fs::Metadata, fs::Metadata,
@@ -18,7 +25,7 @@ use std::{
}; };
use time::OffsetDateTime; use time::OffsetDateTime;
use tokio::fs::{self, File}; use tokio::fs::{self, File};
use tokio::io::ErrorKind; use tokio::io::{AsyncReadExt, AsyncWriteExt, ErrorKind};
use tokio::sync::Mutex; use tokio::sync::Mutex;
use tracing::{debug, warn}; use tracing::{debug, warn};
use uuid::Uuid; use uuid::Uuid;
@@ -26,18 +33,27 @@ use uuid::Uuid;
#[derive(Debug)] #[derive(Debug)]
pub struct FormatInfo { pub struct FormatInfo {
pub id: Option<Uuid>, pub id: Option<Uuid>,
pub _data: Vec<u8>, pub data: Vec<u8>,
pub _file_info: Option<Metadata>, pub file_info: Option<Metadata>,
pub _last_check: Option<OffsetDateTime>, pub last_check: Option<OffsetDateTime>,
} }
impl FormatInfo {} impl FormatInfo {
pub fn last_check_valid(&self) -> bool {
let now = OffsetDateTime::now_utc();
self.file_info.is_some()
&& self.id.is_some()
&& self.last_check.is_some()
&& (now.unix_timestamp() - self.last_check.unwrap().unix_timestamp() <= 1)
}
}
#[derive(Debug)] #[derive(Debug)]
pub struct LocalDisk { pub struct LocalDisk {
pub root: PathBuf, pub root: PathBuf,
pub _format_path: PathBuf, pub format_path: PathBuf,
pub format_info: Mutex<FormatInfo>, pub format_info: Mutex<FormatInfo>,
pub endpoint: Endpoint,
// pub id: Mutex<Option<Uuid>>, // pub id: Mutex<Option<Uuid>>,
// pub format_data: Mutex<Vec<u8>>, // pub format_data: Mutex<Vec<u8>>,
// pub format_file_info: Mutex<Option<Metadata>>, // pub format_file_info: Mutex<Option<Metadata>>,
@@ -79,14 +95,17 @@ impl LocalDisk {
let format_info = FormatInfo { let format_info = FormatInfo {
id, id,
_data: format_data, data: format_data,
_file_info: format_meta, file_info: format_meta,
_last_check: format_last_check, last_check: format_last_check,
}; };
// TODO: DIRECT suport
// TODD: DiskInfo
let disk = Self { let disk = Self {
root, root,
_format_path: format_path, endpoint: ep.clone(),
format_path: format_path,
format_info: Mutex::new(format_info), format_info: Mutex::new(format_info),
// // format_legacy, // // format_legacy,
// format_file_info: Mutex::new(format_meta), // format_file_info: Mutex::new(format_meta),
@@ -99,6 +118,43 @@ impl LocalDisk {
Ok(disk) Ok(disk)
} }
fn is_valid_volname(volname: &str) -> bool {
if volname.len() < 3 {
return false;
}
if cfg!(target_os = "windows") {
// 在 Windows 上,卷名不应该包含保留字符。
// 这个正则表达式匹配了不允许的字符。
if volname.contains('|')
|| volname.contains('<')
|| volname.contains('>')
|| volname.contains('?')
|| volname.contains('*')
|| volname.contains(':')
|| volname.contains('"')
|| volname.contains('\\')
{
return false;
}
} else {
// 对于非 Windows 系统,可能需要其他的验证逻辑。
}
true
}
async fn check_format_json(&self) -> Result<Metadata> {
let md = fs::metadata(&self.format_path).await.map_err(|e| match e.kind() {
ErrorKind::NotFound => DiskError::DiskNotFound,
ErrorKind::PermissionDenied => DiskError::FileAccessDenied,
_ => {
warn!("check_format_json err {:?}", e);
DiskError::CorruptedBackend
}
})?;
Ok(md)
}
async fn make_meta_volumes(&self) -> Result<()> { async fn make_meta_volumes(&self) -> Result<()> {
let buckets = format!("{}/{}", super::RUSTFS_META_BUCKET, super::BUCKET_META_PREFIX); let buckets = format!("{}/{}", super::RUSTFS_META_BUCKET, super::BUCKET_META_PREFIX);
let multipart = format!("{}/{}", super::RUSTFS_META_BUCKET, "multipart"); let multipart = format!("{}/{}", super::RUSTFS_META_BUCKET, "multipart");
@@ -146,10 +202,10 @@ impl LocalDisk {
fs::create_dir_all(dst_data_path.parent().unwrap_or(Path::new("/"))).await?; fs::create_dir_all(dst_data_path.parent().unwrap_or(Path::new("/"))).await?;
} }
debug!( // debug!(
"rename_all from \n {:?} \n to \n {:?} \n skip:{:?}", // "rename_all from \n {:?} \n to \n {:?} \n skip:{:?}",
&src_data_path, &dst_data_path, &skip // &src_data_path, &dst_data_path, &skip
); // );
fs::rename(&src_data_path, &dst_data_path).await?; fs::rename(&src_data_path, &dst_data_path).await?;
@@ -163,7 +219,7 @@ impl LocalDisk {
fs::create_dir_all(parent).await?; fs::create_dir_all(parent).await?;
} }
} }
debug!("move_to_trash from:{:?} to {:?}", &delete_path, &trash_path); // debug!("move_to_trash from:{:?} to {:?}", &delete_path, &trash_path);
// TODO: 清空回收站 // TODO: 清空回收站
if let Err(err) = fs::rename(&delete_path, &trash_path).await { if let Err(err) = fs::rename(&delete_path, &trash_path).await {
match err.kind() { match err.kind() {
@@ -205,15 +261,15 @@ impl LocalDisk {
recursive: bool, recursive: bool,
immediate_purge: bool, immediate_purge: bool,
) -> Result<()> { ) -> Result<()> {
debug!("delete_file {:?}\n base_path:{:?}", &delete_path, &base_path); // debug!("delete_file {:?}\n base_path:{:?}", &delete_path, &base_path);
if is_root_path(base_path) || is_root_path(delete_path) { if is_root_path(base_path) || is_root_path(delete_path) {
debug!("delete_file skip {:?}", &delete_path); // debug!("delete_file skip {:?}", &delete_path);
return Ok(()); return Ok(());
} }
if !delete_path.starts_with(base_path) || base_path == delete_path { if !delete_path.starts_with(base_path) || base_path == delete_path {
debug!("delete_file skip {:?}", &delete_path); // debug!("delete_file skip {:?}", &delete_path);
return Ok(()); return Ok(());
} }
@@ -221,8 +277,9 @@ impl LocalDisk {
self.move_to_trash(delete_path, recursive, immediate_purge).await?; self.move_to_trash(delete_path, recursive, immediate_purge).await?;
} else { } else {
if delete_path.is_dir() { if delete_path.is_dir() {
// debug!("delete_file remove_dir {:?}", &delete_path);
if let Err(err) = fs::remove_dir(&delete_path).await { if let Err(err) = fs::remove_dir(&delete_path).await {
debug!("remove_dir err {:?} when {:?}", &err, &delete_path); // debug!("remove_dir err {:?} when {:?}", &err, &delete_path);
match err.kind() { match err.kind() {
ErrorKind::NotFound => (), ErrorKind::NotFound => (),
// ErrorKind::DirectoryNotEmpty => (), // ErrorKind::DirectoryNotEmpty => (),
@@ -234,9 +291,10 @@ impl LocalDisk {
} }
} }
} }
// debug!("delete_file remove_dir done {:?}", &delete_path);
} else { } else {
if let Err(err) = fs::remove_file(&delete_path).await { if let Err(err) = fs::remove_file(&delete_path).await {
debug!("remove_file err {:?} when {:?}", &err, &delete_path); // debug!("remove_file err {:?} when {:?}", &err, &delete_path);
match err.kind() { match err.kind() {
ErrorKind::NotFound => (), ErrorKind::NotFound => (),
_ => { _ => {
@@ -252,7 +310,7 @@ impl LocalDisk {
Box::pin(self.delete_file(base_path, &PathBuf::from(dir_path), false, false)).await?; Box::pin(self.delete_file(base_path, &PathBuf::from(dir_path), false, false)).await?;
} }
debug!("delete_file done {:?}", &delete_path); // debug!("delete_file done {:?}", &delete_path);
Ok(()) Ok(())
} }
@@ -263,54 +321,125 @@ impl LocalDisk {
volume_dir: impl AsRef<Path>, volume_dir: impl AsRef<Path>,
path: impl AsRef<Path>, path: impl AsRef<Path>,
read_data: bool, read_data: bool,
) -> Result<(Vec<u8>, OffsetDateTime)> { ) -> Result<(Vec<u8>, Option<OffsetDateTime>)> {
let meta_path = path.as_ref().join(Path::new(super::STORAGE_FORMAT_FILE)); let meta_path = path.as_ref().join(Path::new(super::STORAGE_FORMAT_FILE));
if read_data { if read_data {
self.read_all_data(bucket, volume_dir, meta_path).await self.read_all_data_with_dmtime(bucket, volume_dir, meta_path).await
} else { } else {
self.read_metadata_with_dmtime(meta_path).await self.read_all_data_with_dmtime(bucket, volume_dir, meta_path).await
// FIXME: read_metadata only suport
// self.read_metadata_with_dmtime(meta_path).await
} }
} }
async fn read_metadata_with_dmtime(&self, path: impl AsRef<Path>) -> Result<(Vec<u8>, OffsetDateTime)> { // FIXME: read_metadata only suport
let (data, meta) = read_file_all(path).await?; 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())?;
let mut f = utils::fs::open_file(file_path, O_RDONLY).await?;
let meta = f.metadata().await?;
if meta.is_dir() {
return Err(Error::new(DiskError::FileNotFound));
}
let meta = f.metadata().await.map_err(os_err_to_file_err)?;
if meta.is_dir() {
return Err(Error::new(DiskError::FileNotFound));
}
let size = meta.len() as usize;
let mut bytes = Vec::new();
bytes.try_reserve_exact(size)?;
f.read_to_end(&mut bytes).await.map_err(os_err_to_file_err)?;
let modtime = match meta.modified() { let modtime = match meta.modified() {
Ok(md) => OffsetDateTime::from(md), Ok(md) => Some(OffsetDateTime::from(md)),
Err(_) => return Err(Error::msg("Not supported modified on this platform")), Err(_) => None,
}; };
Ok((data, modtime)) Ok((bytes, modtime))
} }
async fn read_all_data( async fn read_all_data(&self, volume: &str, volume_dir: impl AsRef<Path>, file_path: impl AsRef<Path>) -> Result<Vec<u8>> {
&self, // TODO: timeout suport
_bucket: &str, let (data, _) = self.read_all_data_with_dmtime(volume, volume_dir, file_path).await?;
_volume_dir: impl AsRef<Path>, Ok(data)
path: impl AsRef<Path>, }
) -> Result<(Vec<u8>, OffsetDateTime)> {
let (data, meta) = read_file_all(path).await?;
let modtime = match meta.modified() { async fn read_all_data_with_dmtime(
Ok(md) => OffsetDateTime::from(md), &self,
Err(_) => return Err(Error::msg("Not supported modified on this platform")), volume: &str,
volume_dir: impl AsRef<Path>,
file_path: impl AsRef<Path>,
) -> Result<(Vec<u8>, Option<OffsetDateTime>)> {
let mut f = match utils::fs::open_file(file_path.as_ref(), utils::fs::O_RDONLY).await {
Ok(f) => f,
Err(e) => {
if os_is_not_exist(&e) {
if !skip_access_checks(volume) {
if let Err(er) = utils::fs::access(volume_dir.as_ref()).await {
if os_is_not_exist(&er) {
return Err(Error::new(DiskError::VolumeNotFound));
}
}
}
return Err(Error::new(DiskError::FileNotFound));
} else if os_is_permission(&e) {
return Err(Error::new(DiskError::FileAccessDenied));
} else if is_sys_err_not_dir(&e) || is_sys_err_is_dir(&e) {
return Err(Error::new(DiskError::FileNotFound));
} else if is_sys_err_handle_invalid(&e) {
return Err(Error::new(DiskError::FileNotFound));
} else if is_sys_err_io(&e) {
return Err(Error::new(DiskError::FaultyDisk));
} else if is_sys_err_too_many_files(&e) {
return Err(Error::new(DiskError::TooManyOpenFiles));
} else if is_sys_err_invalid_arg(&e) {
if let Ok(meta) = utils::fs::lstat(file_path.as_ref()).await {
if meta.is_dir() {
return Err(Error::new(DiskError::FileNotFound));
}
}
return Err(Error::new(DiskError::UnsupportedDisk));
}
return Err(Error::new(e));
}
}; };
Ok((data, modtime)) let meta = f.metadata().await.map_err(os_err_to_file_err)?;
if meta.is_dir() {
return Err(Error::new(DiskError::FileNotFound));
}
let size = meta.len() as usize;
let mut bytes = Vec::new();
bytes.try_reserve_exact(size)?;
f.read_to_end(&mut bytes).await.map_err(os_err_to_file_err)?;
let modtime = match meta.modified() {
Ok(md) => Some(OffsetDateTime::from(md)),
Err(_) => None,
};
Ok((bytes, modtime))
} }
async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &Vec<FileInfo>) -> Result<()> { async fn delete_versions_internal(&self, volume: &str, path: &str, fis: &Vec<FileInfo>) -> Result<()> {
let volume_dir = self.get_bucket_path(volume)?; let volume_dir = self.get_bucket_path(volume)?;
let xlpath = self.get_object_path(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str())?; let xlpath = self.get_object_path(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str())?;
let (data, _) = match self.read_all_data(volume, volume_dir.as_path(), &xlpath).await { let data = self
Ok(res) => res, .read_all_data(volume, volume_dir.as_path(), &xlpath)
Err(_err) => { .await
// TODO: check if not found return err .unwrap_or_default();
(Vec::new(), OffsetDateTime::UNIX_EPOCH)
}
};
if data.is_empty() { if data.is_empty() {
return Err(Error::new(DiskError::FileNotFound)); return Err(Error::new(DiskError::FileNotFound));
@@ -340,11 +469,114 @@ impl LocalDisk {
// 更新xl.meta // 更新xl.meta
let buf = fm.marshal_msg()?; let buf = fm.marshal_msg()?;
self.write_all(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str(), buf) let volume_dir = self.get_bucket_path(&volume)?;
.await?;
self.write_all_private(
volume,
format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str(),
&buf,
true,
volume_dir,
)
.await?;
Ok(()) Ok(())
} }
async fn write_all_meta(&self, volume: &str, path: &str, buf: &[u8], sync: bool) -> Result<()> {
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())?;
let tmp_volume_dir = self.get_bucket_path(super::RUSTFS_META_TMP_BUCKET)?;
let tmp_file_path = tmp_volume_dir.join(Path::new(Uuid::new_v4().to_string().as_str()));
self.write_all_internal(&tmp_file_path, buf, sync, tmp_volume_dir).await?;
os::rename_all(tmp_file_path, file_path, volume_dir).await
}
// write_all_public for trail
async fn write_all_public(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()> {
if volume == super::RUSTFS_META_BUCKET && path == super::FORMAT_CONFIG_FILE {
let mut format_info = self.format_info.lock().await;
format_info.data = data.clone();
}
let volume_dir = self.get_bucket_path(&volume)?;
self.write_all_private(&volume, &path, &data, true, volume_dir).await?;
Ok(())
}
// write_all_private with check_path_length
pub async fn write_all_private(
&self,
volume: &str,
path: &str,
buf: &[u8],
sync: bool,
skip_parent: impl AsRef<Path>,
) -> Result<()> {
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())?;
self.write_all_internal(file_path, buf, sync, skip_parent).await
}
// write_all_internal do write file
pub async fn write_all_internal(
&self,
file_path: impl AsRef<Path>,
data: impl AsRef<[u8]>,
sync: bool,
skip_parent: impl AsRef<Path>,
) -> Result<()> {
let flags = utils::fs::O_CREATE | utils::fs::O_WRONLY | utils::fs::O_TRUNC;
let mut f = {
if sync {
// TODO: suport sync
self.open_file(file_path.as_ref(), flags, skip_parent.as_ref()).await?
} else {
self.open_file(file_path.as_ref(), flags, skip_parent.as_ref()).await?
}
};
f.write_all(data.as_ref()).await?;
Ok(())
}
async fn open_file(&self, path: impl AsRef<Path>, mode: usize, skip_parent: impl AsRef<Path>) -> Result<File> {
let mut skip_parent = skip_parent.as_ref();
if skip_parent.as_os_str().is_empty() {
skip_parent = self.root.as_path();
}
if let Some(parent) = path.as_ref().parent() {
os::make_dir_all(parent, skip_parent).await?;
}
let f = utils::fs::open_file(path.as_ref(), mode).await.map_err(|e| {
if is_sys_err_io(&e) {
Error::new(DiskError::IsNotRegular)
} 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)
}
})?;
Ok(f)
}
} }
fn is_root_path(path: impl AsRef<Path>) -> bool { fn is_root_path(path: impl AsRef<Path>) -> bool {
@@ -373,14 +605,6 @@ pub async fn read_file_exists(path: impl AsRef<Path>) -> Result<(Vec<u8>, Option
Ok((data, meta)) Ok((data, meta))
} }
pub async fn write_all_internal(p: impl AsRef<Path>, data: impl AsRef<[u8]>) -> Result<()> {
// create top dir if not exists
fs::create_dir_all(&p.as_ref().parent().unwrap_or_else(|| Path::new("."))).await?;
fs::write(&p, data).await?;
Ok(())
}
pub async fn read_file_all(path: impl AsRef<Path>) -> Result<(Vec<u8>, Metadata)> { pub async fn read_file_all(path: impl AsRef<Path>) -> Result<(Vec<u8>, Metadata)> {
let p = path.as_ref(); let p = path.as_ref();
let meta = read_file_metadata(&path).await?; let meta = read_file_metadata(&path).await?;
@@ -428,9 +652,23 @@ fn skip_access_checks(p: impl AsRef<str>) -> bool {
#[async_trait::async_trait] #[async_trait::async_trait]
impl DiskAPI for LocalDisk { impl DiskAPI for LocalDisk {
fn to_string(&self) -> String {
self.root.to_string_lossy().to_string()
}
fn is_local(&self) -> bool { fn is_local(&self) -> bool {
true true
} }
fn host_name(&self) -> String {
self.endpoint.host_port()
}
async fn is_online(&self) -> bool {
true
}
fn endpoint(&self) -> Endpoint {
self.endpoint.clone()
}
async fn close(&self) -> Result<()> { async fn close(&self) -> Result<()> {
Ok(()) Ok(())
} }
@@ -438,12 +676,67 @@ impl DiskAPI for LocalDisk {
self.root.clone() self.root.clone()
} }
async fn get_disk_id(&self) -> Option<Uuid> { fn get_disk_location(&self) -> DiskLocation {
warn!("local get_disk_id"); DiskLocation {
// TODO: check format file pool_idx: self.endpoint.pool_idx,
let format_info = self.format_info.lock().await; set_idx: self.endpoint.set_idx,
disk_idx: self.endpoint.pool_idx,
}
}
format_info.id.clone() async fn get_disk_id(&self) -> Result<Option<Uuid>> {
// TODO: check format file
let mut format_info = self.format_info.lock().await;
let id = format_info.id.clone();
if format_info.last_check_valid() {
return Ok(id);
}
let file_meta = self.check_format_json().await?;
if let Some(file_info) = &format_info.file_info {
if utils::fs::same_file(&file_meta, file_info) {
format_info.last_check = Some(OffsetDateTime::now_utc());
return Ok(id);
}
}
let b = fs::read(&self.format_path).await.map_err(|e| match e.kind() {
ErrorKind::NotFound => DiskError::DiskNotFound,
ErrorKind::PermissionDenied => DiskError::FileAccessDenied,
_ => {
warn!("check_format_json err {:?}", e);
DiskError::CorruptedBackend
}
})?;
let fm = FormatV3::try_from(b.as_slice()).map_err(|e| {
warn!("decode format.json err {:?}", e);
DiskError::CorruptedBackend
})?;
let (m, n) = fm.find_disk_index_by_disk_id(fm.erasure.this)?;
let disk_id = fm.erasure.this;
match (self.endpoint.set_idx, self.endpoint.disk_idx) {
(Some(set_idx), Some(disk_idx)) => {
if m != set_idx || n != disk_idx {
return Err(Error::new(DiskError::InconsistentDisk));
}
}
_ => return Err(Error::new(DiskError::InconsistentDisk)),
}
format_info.id = Some(disk_id);
format_info.file_info = Some(file_meta);
format_info.data = b;
format_info.last_check = Some(OffsetDateTime::now_utc());
Ok(Some(disk_id))
// TODO: 判断源文件id,是否有效 // TODO: 判断源文件id,是否有效
} }
@@ -455,19 +748,16 @@ impl DiskAPI for LocalDisk {
} }
#[must_use] #[must_use]
async fn read_all(&self, volume: &str, path: &str) -> Result<Bytes> { async fn read_all(&self, volume: &str, path: &str) -> Result<Vec<u8>> {
// TOFIX:
let p = self.get_object_path(volume, path)?; let p = self.get_object_path(volume, path)?;
let (data, _) = read_file_all(&p).await?; let (data, _) = read_file_all(&p).await?;
Ok(Bytes::from(data)) Ok(data)
} }
async fn write_all(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()> { async fn write_all(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()> {
let p = self.get_object_path(volume, path)?; self.write_all_public(volume, path, data).await
write_all_internal(p, data).await?;
Ok(())
} }
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> { async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()> {
@@ -494,15 +784,95 @@ impl DiskAPI for LocalDisk {
Ok(()) Ok(())
} }
async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec<u8>) -> Result<()> {
async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> { let src_volume_dir = self.get_bucket_path(src_volume)?;
let src_volume_path = self.get_bucket_path(src_volume)?; let dst_volume_dir = self.get_bucket_path(dst_volume)?;
if !skip_access_checks(src_volume) { if !skip_access_checks(src_volume) {
check_volume_exists(&src_volume_path).await?; utils::fs::access(&src_volume_dir).await.map_err(map_err_not_exists)?
} }
if !skip_access_checks(dst_volume) { if !skip_access_checks(dst_volume) {
let vol_path = self.get_bucket_path(dst_volume)?; utils::fs::access(&dst_volume_dir).await.map_err(map_err_not_exists)?
check_volume_exists(&vol_path).await?; }
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));
}
let src_file_path = src_volume_dir.join(Path::new(src_path));
let dst_file_path = dst_volume_dir.join(Path::new(dst_path));
check_path_length(&src_file_path.to_string_lossy().to_string())?;
check_path_length(&dst_file_path.to_string_lossy().to_string())?;
if src_is_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));
}
if !os_is_not_exist(&e) {
return Err(Error::new(e));
}
None
}
};
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 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 Err(err) = self.write_all(&dst_volume, format!("{}.meta", dst_path).as_str(), meta).await {
if let Some(e) = err.to_io_err() {
return Err(os_err_to_file_err(e));
}
return Err(err);
}
if let Some(parent) = src_file_path.parent() {
self.delete_file(&src_volume_dir, &parent.to_path_buf(), false, false).await?;
}
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)?;
if !skip_access_checks(src_volume) {
utils::fs::access(&src_volume_path).await.map_err(map_err_not_exists)?;
}
if !skip_access_checks(dst_volume) {
utils::fs::access(&dst_volume_path).await.map_err(map_err_not_exists)?;
} }
let srcp = self.get_object_path(src_volume, src_path)?; let srcp = self.get_object_path(src_volume, src_path)?;
@@ -543,6 +913,13 @@ impl DiskAPI for LocalDisk {
} }
async fn create_file(&self, _origvolume: &str, volume: &str, path: &str, _file_size: usize) -> Result<FileWriter> { 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)?; let fpath = self.get_object_path(volume, path)?;
debug!("CreateFile fpath: {:?}", fpath); debug!("CreateFile fpath: {:?}", fpath);
@@ -595,7 +972,7 @@ impl DiskAPI for LocalDisk {
async fn read_file(&self, volume: &str, path: &str) -> Result<FileReader> { async fn read_file(&self, volume: &str, path: &str) -> Result<FileReader> {
let p = self.get_object_path(volume, path)?; let p = self.get_object_path(volume, path)?;
debug!("read_file {:?}", &p); // debug!("read_file {:?}", &p);
let file = File::options().read(true).open(&p).await?; let file = File::options().read(true).open(&p).await?;
Ok(FileReader::Local(LocalFileReader::new(file))) Ok(FileReader::Local(LocalFileReader::new(file)))
@@ -671,10 +1048,7 @@ impl DiskAPI for LocalDisk {
let (fdata, _) = match self.read_metadata_with_dmtime(&fpath).await { let (fdata, _) = match self.read_metadata_with_dmtime(&fpath).await {
Ok(res) => res, Ok(res) => res,
Err(_) => { Err(_) => (Vec::new(), None),
// TODO: check err
(Vec::new(), OffsetDateTime::UNIX_EPOCH)
}
}; };
meta.metadata = fdata; meta.metadata = fdata;
@@ -740,8 +1114,6 @@ impl DiskAPI for LocalDisk {
} }
} }
warn!("get meta {:?}", &meta);
let mut skip_parent = dst_volume_path.clone(); let mut skip_parent = dst_volume_path.clone();
if !dst_buf.is_empty() { if !dst_buf.is_empty() {
skip_parent = PathBuf::from(&dst_file_path.parent().unwrap_or(Path::new("/"))); skip_parent = PathBuf::from(&dst_file_path.parent().unwrap_or(Path::new("/")));
@@ -763,14 +1135,15 @@ impl DiskAPI for LocalDisk {
let fm_data = meta.marshal_msg()?; let fm_data = meta.marshal_msg()?;
// 写入xl.meta // 写入xl.meta
write_all_internal(&src_file_path, fm_data).await?; 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; let no_inline = src_data_path.has_root() && fi.data.is_none() && fi.size > 0;
if no_inline { if no_inline {
self.rename_all(&src_data_path, &dst_data_path, &skip_parent).await?; self.rename_all(&src_data_path, &dst_data_path, &skip_parent).await?;
} }
warn!("old_data_dir {:?}", old_data_dir); // warn!("old_data_dir {:?}", old_data_dir);
// 有旧目录,把old xl.meta存到旧目录里 // 有旧目录,把old xl.meta存到旧目录里
if old_data_dir.is_some() { if old_data_dir.is_some() {
self.write_all( self.write_all(
@@ -794,60 +1167,60 @@ impl DiskAPI for LocalDisk {
.await?; .await?;
} }
Ok(RenameDataResp { old_data_dir }) Ok(RenameDataResp {
old_data_dir: old_data_dir,
sign: None, // TODO:
})
} }
async fn make_volumes(&self, volumes: Vec<&str>) -> Result<()> { async fn make_volumes(&self, volumes: Vec<&str>) -> Result<()> {
for vol in volumes { for vol in volumes {
if let Err(e) = self.make_volume(vol).await { if let Err(e) = self.make_volume(vol).await {
match &e.downcast_ref::<DiskError>() { if !DiskError::VolumeExists.is(&e) {
Some(DiskError::VolumeExists) => Ok(()), return Err(e);
Some(_) => Err(e), }
None => Err(e),
}?;
} }
// TODO: health check // TODO: health check
} }
Ok(()) Ok(())
} }
async fn make_volume(&self, volume: &str) -> Result<()> { async fn make_volume(&self, volume: &str) -> Result<()> {
if !Self::is_valid_volname(volume) {
return Err(Error::msg("Invalid arguments specified"));
}
let p = self.get_bucket_path(volume)?; let p = self.get_bucket_path(volume)?;
match File::open(&p).await {
Ok(_) => (), if let Err(e) = utils::fs::access(&p).await {
Err(e) => match e.kind() { if os_is_not_exist(&e) {
ErrorKind::NotFound => { os::make_dir_all(&p, self.root.as_path()).await?;
fs::create_dir_all(&p).await?; }
return Ok(());
}
_ => return Err(Error::from(e)),
},
} }
Err(Error::from(DiskError::VolumeExists)) Err(Error::from(DiskError::VolumeExists))
} }
async fn list_volumes(&self) -> Result<Vec<VolumeInfo>> { async fn list_volumes(&self) -> Result<Vec<VolumeInfo>> {
let mut entries = fs::read_dir(&self.root).await?;
let mut volumes = Vec::new(); let mut volumes = Vec::new();
while let Some(entry) = entries.next_entry().await? { let entries = os::read_dir(&self.root, 0).await.map_err(|e| {
if let Ok(metadata) = entry.metadata().await { if DiskError::FileAccessDenied.is(&e) {
// if !metadata.is_dir() { Error::new(DiskError::DiskAccessDenied)
// continue; } else if DiskError::FileNotFound.is(&e) {
// } Error::new(DiskError::DiskAccessDenied)
} else {
let name = entry.file_name().to_string_lossy().to_string(); e
let created = match metadata.created() {
Ok(md) => Some(OffsetDateTime::from(md)),
Err(_) => {
warn!("Not supported created on this platform");
None
}
};
volumes.push(VolumeInfo { name, created });
} }
})?;
for entry in entries {
if !utils::path::has_suffix(&entry, SLASH_SEPARATOR) || !Self::is_valid_volname(&entry) {
continue;
}
volumes.push(VolumeInfo {
name: entry,
created: None,
});
} }
Ok(volumes) Ok(volumes)
@@ -869,7 +1242,64 @@ impl DiskAPI for LocalDisk {
created: modtime, created: modtime,
}) })
} }
async fn delete_paths(&self, volume: &str, paths: &[&str]) -> Result<()> {
let volume_dir = self.get_bucket_path(volume)?;
if !skip_access_checks(volume) {
utils::fs::access(&volume_dir)
.await
.map_err(|e| convert_access_error(e, DiskError::VolumeAccessDenied))?
}
for path in paths.iter() {
let file_path = volume_dir.join(Path::new(path));
check_path_length(&file_path.to_string_lossy().to_string())?;
self.move_to_trash(&file_path, false, false).await?;
}
Ok(())
}
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: UpdateMetadataOpts) -> Result<()> {
if let Some(_) = &fi.metadata {
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())?;
let buf = self
.read_all(&volume, format!("{}/{}", &path, super::STORAGE_FORMAT_FILE).as_str())
.await
.map_err(|e| {
if DiskError::FileNotFound.is(&e) && fi.version_id.is_some() {
Error::new(DiskError::FileVersionNotFound)
} else {
e
}
})?;
if !FileMeta::is_xl_format(buf.as_slice()) {
return Err(Error::new(DiskError::FileVersionNotFound));
}
let mut xl_meta = FileMeta::load(buf.as_slice())?;
xl_meta.update_object_version(fi)?;
let wbuf = xl_meta.marshal_msg()?;
return self
.write_all_meta(
&volume,
format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str(),
&wbuf,
!opts.no_persistence,
)
.await;
}
Err(Error::msg("Invalid Argument"))
}
async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
let p = self.get_object_path(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str())?; let p = self.get_object_path(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str())?;
@@ -890,7 +1320,8 @@ impl DiskAPI for LocalDisk {
let fm_data = meta.marshal_msg()?; let fm_data = meta.marshal_msg()?;
write_all_internal(p, fm_data).await?; self.write_all(volume, format!("{}/{}", path, super::STORAGE_FORMAT_FILE).as_str(), fm_data)
.await?;
return Ok(()); return Ok(());
} }
@@ -924,6 +1355,20 @@ impl DiskAPI for LocalDisk {
Ok(RawFileInfo { buf }) Ok(RawFileInfo { buf })
} }
async fn delete_version(
&self,
volume: &str,
_path: &str,
_fi: FileInfo,
_force_del_marker: bool,
_opts: DeleteOptions,
) -> Result<RawFileInfo> {
let _volume_dir = self.get_bucket_path(&volume)?;
// self.read_all_data(bucket, volume_dir, path);
unimplemented!()
}
async fn delete_versions( async fn delete_versions(
&self, &self,
volume: &str, volume: &str,
+87 -32
View File
@@ -2,6 +2,7 @@ pub mod endpoint;
pub mod error; pub mod error;
pub mod format; pub mod format;
mod local; mod local;
pub mod os;
mod remote; mod remote;
pub const RUSTFS_META_BUCKET: &str = ".rustfs.sys"; pub const RUSTFS_META_BUCKET: &str = ".rustfs.sys";
@@ -18,7 +19,8 @@ use crate::{
file_meta::FileMeta, file_meta::FileMeta,
store_api::{FileInfo, RawFileInfo}, store_api::{FileInfo, RawFileInfo},
}; };
use bytes::Bytes;
use endpoint::Endpoint;
use protos::proto_gen::node_service::{node_service_client::NodeServiceClient, ReadAtRequest, WriteRequest}; use protos::proto_gen::node_service::{node_service_client::NodeServiceClient, ReadAtRequest, WriteRequest};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::{fmt::Debug, io::SeekFrom, path::PathBuf, sync::Arc}; use std::{fmt::Debug, io::SeekFrom, path::PathBuf, sync::Arc};
@@ -44,39 +46,51 @@ pub async fn new_disk(ep: &endpoint::Endpoint, opt: &DiskOption) -> Result<DiskS
#[async_trait::async_trait] #[async_trait::async_trait]
pub trait DiskAPI: Debug + Send + Sync + 'static { pub trait DiskAPI: Debug + Send + Sync + 'static {
fn to_string(&self) -> String;
async fn is_online(&self) -> bool;
fn is_local(&self) -> bool; fn is_local(&self) -> bool;
fn path(&self) -> PathBuf; // LastConn
fn host_name(&self) -> String;
fn endpoint(&self) -> Endpoint;
async fn close(&self) -> Result<()>; async fn close(&self) -> Result<()>;
async fn get_disk_id(&self) -> Option<Uuid>; async fn get_disk_id(&self) -> Result<Option<Uuid>>;
async fn set_disk_id(&self, id: Option<Uuid>) -> Result<()>; async fn set_disk_id(&self, id: Option<Uuid>) -> Result<()>;
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()>; fn path(&self) -> PathBuf;
async fn read_all(&self, volume: &str, path: &str) -> Result<Bytes>; fn get_disk_location(&self) -> DiskLocation;
async fn write_all(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()>;
async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()>; // Healing
async fn create_file(&self, origvolume: &str, volume: &str, path: &str, file_size: usize) -> Result<FileWriter>; // DiskInfo
async fn append_file(&self, volume: &str, path: &str) -> Result<FileWriter>; // NSScanner
async fn read_file(&self, volume: &str, path: &str) -> Result<FileReader>;
// 读目录下的所有文件、目录 // Volume operations.
async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> Result<Vec<String>>; async fn make_volume(&self, volume: &str) -> Result<()>;
async fn make_volumes(&self, volume: Vec<&str>) -> Result<()>;
async fn list_volumes(&self) -> Result<Vec<VolumeInfo>>;
async fn stat_volume(&self, volume: &str) -> Result<VolumeInfo>;
async fn delete_volume(&self, volume: &str) -> Result<()>;
// 并发边读边写 TODO: wr io.Writer // 并发边读边写 TODO: wr io.Writer
async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>>; async fn walk_dir(&self, opts: WalkDirOptions) -> Result<Vec<MetaCacheEntry>>;
async fn rename_data(
// Metadata operations
async fn delete_version(
&self, &self,
src_volume: &str, volume: &str,
src_path: &str, path: &str,
file_info: FileInfo, fi: FileInfo,
dst_volume: &str, force_del_marker: bool,
dst_path: &str, opts: DeleteOptions,
) -> Result<RenameDataResp>; ) -> Result<RawFileInfo>;
async fn delete_versions(
async fn make_volumes(&self, volume: Vec<&str>) -> Result<()>; &self,
async fn delete_volume(&self, volume: &str) -> Result<()>; volume: &str,
async fn list_volumes(&self) -> Result<Vec<VolumeInfo>>; versions: Vec<FileInfoVersions>,
async fn make_volume(&self, volume: &str) -> Result<()>; opts: DeleteOptions,
async fn stat_volume(&self, volume: &str) -> Result<VolumeInfo>; ) -> Result<Vec<Option<Error>>>;
async fn delete_paths(&self, volume: &str, paths: &[&str]) -> Result<()>;
async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()>; async fn write_metadata(&self, org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()>;
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: UpdateMetadataOpts) -> Result<()>;
async fn read_version( async fn read_version(
&self, &self,
org_volume: &str, org_volume: &str,
@@ -86,13 +100,50 @@ pub trait DiskAPI: Debug + Send + Sync + 'static {
opts: &ReadOptions, opts: &ReadOptions,
) -> Result<FileInfo>; ) -> Result<FileInfo>;
async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result<RawFileInfo>; async fn read_xl(&self, volume: &str, path: &str, read_data: bool) -> Result<RawFileInfo>;
async fn delete_versions( async fn rename_data(
&self, &self,
volume: &str, src_volume: &str,
versions: Vec<FileInfoVersions>, src_path: &str,
opts: DeleteOptions, file_info: FileInfo,
) -> Result<Vec<Option<Error>>>; dst_volume: &str,
dst_path: &str,
) -> Result<RenameDataResp>;
// File operations.
// 读目录下的所有文件、目录
async fn list_dir(&self, origvolume: &str, volume: &str, dir_path: &str, count: i32) -> Result<Vec<String>>;
async fn read_file(&self, volume: &str, path: &str) -> Result<FileReader>;
async fn append_file(&self, volume: &str, path: &str) -> Result<FileWriter>;
async fn create_file(&self, origvolume: &str, volume: &str, path: &str, file_size: usize) -> Result<FileWriter>;
// ReadFileStream
async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()>;
async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec<u8>) -> Result<()>;
// CheckParts
async fn delete(&self, volume: &str, path: &str, opt: DeleteOptions) -> Result<()>;
// VerifyFile
// StatInfoFile
// ReadParts
async fn read_multiple(&self, req: ReadMultipleReq) -> Result<Vec<ReadMultipleResp>>; async fn read_multiple(&self, req: ReadMultipleReq) -> Result<Vec<ReadMultipleResp>>;
// CleanAbandonedData
async fn write_all(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()>;
async fn read_all(&self, volume: &str, path: &str) -> Result<Vec<u8>>;
}
#[derive(Debug, Serialize, Deserialize)]
pub struct UpdateMetadataOpts {
pub no_persistence: bool,
}
pub struct DiskLocation {
pub pool_idx: Option<usize>,
pub set_idx: Option<usize>,
pub disk_idx: Option<usize>,
}
impl DiskLocation {
pub fn valid(&self) -> bool {
self.pool_idx.is_some() && self.set_idx.is_some() && self.disk_idx.is_some()
}
} }
#[derive(Debug, Default, Clone, Serialize, Deserialize)] #[derive(Debug, Default, Clone, Serialize, Deserialize)]
@@ -208,20 +259,24 @@ impl MetaCacheEntry {
} }
} }
#[derive(Debug, Default)]
pub struct DiskOption { pub struct DiskOption {
pub cleanup: bool, pub cleanup: bool,
pub health_check: bool, pub health_check: bool,
} }
#[derive(Serialize, Deserialize)] #[derive(Debug, Default, Serialize, Deserialize)]
pub struct RenameDataResp { pub struct RenameDataResp {
pub old_data_dir: Option<Uuid>, pub old_data_dir: Option<Uuid>,
pub sign: Option<Vec<u8>>,
} }
#[derive(Debug, Clone, Default, Serialize, Deserialize)] #[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct DeleteOptions { pub struct DeleteOptions {
pub recursive: bool, pub recursive: bool,
pub immediate: bool, pub immediate: bool,
pub undo_write: bool,
pub old_data_dir: Option<Uuid>,
} }
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
+219
View File
@@ -0,0 +1,219 @@
use std::{
io,
path::{Component, Path},
};
use tokio::fs;
use crate::{
disk::error::{is_sys_err_not_dir, is_sys_err_path_not_found, os_is_not_exist},
error::{Error, Result},
utils,
};
use super::error::{os_err_to_file_err, os_is_exist, DiskError};
pub fn check_path_length(path_name: &str) -> Result<()> {
// Apple OS X path length is limited to 1016
if cfg!(target_os = "macos") && path_name.len() > 1016 {
return Err(Error::new(DiskError::FileNameTooLong));
}
// Disallow more than 1024 characters on windows, there
// are no known name_max limits on Windows.
if cfg!(target_os = "windows") && path_name.len() > 1024 {
return Err(Error::new(DiskError::FileNameTooLong));
}
// On Unix we reject paths if they are just '.', '..' or '/'
let invalid_paths = [".", "..", "/"];
if invalid_paths.contains(&path_name) {
return Err(Error::new(DiskError::FileAccessDenied));
}
// Check each path segment length is > 255 on all Unix
// platforms, look for this value as NAME_MAX in
// /usr/include/linux/limits.h
let mut count = 0usize;
for c in path_name.chars() {
match c {
'/' | '\\' if cfg!(target_os = "windows") => count = 0, // Reset
_ => {
count += 1;
if count > 255 {
return Err(Error::new(DiskError::FileNameTooLong));
}
}
}
}
// Success.
Ok(())
}
pub async fn make_dir_all(path: impl AsRef<Path>, base_dir: impl AsRef<Path>) -> Result<()> {
check_path_length(path.as_ref().to_string_lossy().to_string().as_str())?;
if let Err(e) = reliable_mkdir_all(path.as_ref(), base_dir.as_ref()).await {
if is_sys_err_not_dir(&e) {
return Err(Error::new(DiskError::FileAccessDenied));
} else if is_sys_err_path_not_found(&e) {
return Err(Error::new(DiskError::FileAccessDenied));
}
return Err(os_err_to_file_err(e));
}
Ok(())
}
// read_dir count read limit. when count == 0 unlimit.
pub async fn read_dir(path: impl AsRef<Path>, count: usize) -> 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
}
};
while let Some(entry) = entries.next_entry().await? {
let name = entry.file_name().to_string_lossy().to_string();
if name == "" || name == "." || name == ".." {
continue;
}
let file_type = entry.file_type().await?;
if file_type.is_dir() {
count -= 1;
volumes.push(format!("{}{}", name, utils::path::SLASH_SEPARATOR));
if count == 0 {
break;
}
}
}
Ok(volumes)
}
pub async fn rename_all(
src_file_path: impl AsRef<Path>,
dst_file_path: impl AsRef<Path>,
base_dir: impl AsRef<Path>,
) -> Result<()> {
reliable_rename(src_file_path, dst_file_path, base_dir).await.map_err(|e| {
if is_sys_err_not_dir(&e) || !os_is_not_exist(&e) {
Error::new(DiskError::FileAccessDenied)
} else if is_sys_err_path_not_found(&e) {
Error::new(DiskError::FileAccessDenied)
} else if os_is_not_exist(&e) {
Error::new(DiskError::FileNotFound)
} else if os_is_exist(&e) {
Error::new(DiskError::IsNotRegular)
} else {
Error::new(e)
}
})?;
Ok(())
}
pub async fn reliable_rename(
src_file_path: impl AsRef<Path>,
dst_file_path: impl AsRef<Path>,
base_dir: impl AsRef<Path>,
) -> io::Result<()> {
if let Some(parent) = dst_file_path.as_ref().parent() {
reliable_mkdir_all(parent, base_dir.as_ref()).await?;
}
let mut i = 0;
loop {
if let Err(e) = utils::fs::rename(src_file_path.as_ref(), dst_file_path.as_ref()).await {
if os_is_not_exist(&e) && i == 0 {
i += 1;
continue;
}
return Err(e);
}
break;
}
Ok(())
}
pub async fn reliable_mkdir_all(path: impl AsRef<Path>, base_dir: impl AsRef<Path>) -> io::Result<()> {
let mut i = 0;
let mut base_dir = base_dir.as_ref();
loop {
if let Err(e) = os_mkdir_all(path.as_ref(), base_dir).await {
if os_is_not_exist(&e) && i == 0 {
i += 1;
if let Some(base_parent) = base_dir.parent() {
if let Some(c) = base_parent.components().next() {
if c != Component::RootDir {
base_dir = base_parent
}
}
}
continue;
}
return Err(e);
}
break;
}
Ok(())
}
pub async fn os_mkdir_all(dir_path: impl AsRef<Path>, base_dir: impl AsRef<Path>) -> io::Result<()> {
if !base_dir.as_ref().to_string_lossy().is_empty() {
if base_dir.as_ref().starts_with(dir_path.as_ref()) {
return Ok(());
}
}
if let Some(parent) = dir_path.as_ref().parent() {
// 不支持递归,直接create_dir_all了
if let Err(e) = utils::fs::make_dir_all(&parent).await {
if os_is_exist(&e) {
return Ok(());
}
return Err(e);
}
// Box::pin(os_mkdir_all(&parent, &base_dir)).await?;
}
if let Err(e) = utils::fs::mkdir(dir_path.as_ref()).await {
if os_is_exist(&e) {
return Ok(());
}
return Err(e);
}
Ok(())
}
#[cfg(test)]
mod test {
use super::*;
#[tokio::test]
async fn test_make_dir() {}
}
+134 -16
View File
@@ -1,14 +1,13 @@
use std::{path::PathBuf, sync::Arc, time::Duration}; use std::{path::PathBuf, sync::Arc, time::Duration};
use bytes::Bytes;
use futures::lock::Mutex; use futures::lock::Mutex;
use protos::{ use protos::{
node_service_time_out_client, node_service_time_out_client,
proto_gen::node_service::{ proto_gen::node_service::{
node_service_client::NodeServiceClient, DeleteRequest, DeleteVersionsRequest, DeleteVolumeRequest, ListDirRequest, node_service_client::NodeServiceClient, DeletePathsRequest, DeleteRequest, DeleteVersionRequest, DeleteVersionsRequest,
ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, ReadAllRequest, ReadMultipleRequest, ReadVersionRequest, DeleteVolumeRequest, ListDirRequest, ListVolumesRequest, MakeVolumeRequest, MakeVolumesRequest, ReadAllRequest,
ReadXlRequest, RenameDataRequest, RenameFileRequst, StatVolumeRequest, WalkDirRequest, WriteAllRequest, ReadMultipleRequest, ReadVersionRequest, ReadXlRequest, RenameDataRequest, RenameFileRequst, RenamePartRequst,
WriteMetadataRequest, StatVolumeRequest, UpdateMetadataRequest, WalkDirRequest, WriteAllRequest, WriteMetadataRequest,
}, },
DEFAULT_GRPC_SERVER_MESSAGE_LEN, DEFAULT_GRPC_SERVER_MESSAGE_LEN,
}; };
@@ -28,9 +27,9 @@ use crate::{
}; };
use super::{ use super::{
endpoint::Endpoint, DeleteOptions, DiskAPI, DiskOption, FileInfoVersions, FileReader, FileWriter, MetaCacheEntry, endpoint::Endpoint, DeleteOptions, DiskAPI, DiskLocation, DiskOption, FileInfoVersions, FileReader, FileWriter,
ReadMultipleReq, ReadMultipleResp, ReadOptions, RemoteFileReader, RemoteFileWriter, RenameDataResp, VolumeInfo, MetaCacheEntry, ReadMultipleReq, ReadMultipleResp, ReadOptions, RemoteFileReader, RemoteFileWriter, RenameDataResp,
WalkDirOptions, UpdateMetadataOpts, VolumeInfo, WalkDirOptions,
}; };
#[derive(Debug)] #[derive(Debug)]
@@ -39,6 +38,7 @@ pub struct RemoteDisk {
channel: Arc<RwLock<Option<Channel>>>, channel: Arc<RwLock<Option<Channel>>>,
url: url::Url, url: url::Url,
pub root: PathBuf, pub root: PathBuf,
endpoint: Endpoint,
} }
impl RemoteDisk { impl RemoteDisk {
@@ -50,6 +50,7 @@ impl RemoteDisk {
url: ep.url.clone(), url: ep.url.clone(),
root, root,
id: Mutex::new(None), id: Mutex::new(None),
endpoint: ep.clone(),
}) })
} }
@@ -95,9 +96,27 @@ impl RemoteDisk {
// TODO: all api need to handle errors // TODO: all api need to handle errors
#[async_trait::async_trait] #[async_trait::async_trait]
impl DiskAPI for RemoteDisk { impl DiskAPI for RemoteDisk {
fn to_string(&self) -> String {
self.endpoint.to_string()
}
fn is_local(&self) -> bool { fn is_local(&self) -> bool {
false false
} }
fn host_name(&self) -> String {
self.endpoint.host_port()
}
async fn is_online(&self) -> bool {
// TODO: 连接状态
if let Ok(_) = self.get_client_v2().await {
return true;
}
false
}
fn endpoint(&self) -> Endpoint {
self.endpoint.clone()
}
async fn close(&self) -> Result<()> { async fn close(&self) -> Result<()> {
Ok(()) Ok(())
} }
@@ -105,8 +124,16 @@ impl DiskAPI for RemoteDisk {
self.root.clone() self.root.clone()
} }
async fn get_disk_id(&self) -> Option<Uuid> { fn get_disk_location(&self) -> DiskLocation {
self.id.lock().await.clone() DiskLocation {
pool_idx: self.endpoint.pool_idx,
set_idx: self.endpoint.set_idx,
disk_idx: self.endpoint.pool_idx,
}
}
async fn get_disk_id(&self) -> Result<Option<Uuid>> {
Ok(self.id.lock().await.clone())
} }
async fn set_disk_id(&self, id: Option<Uuid>) -> Result<()> { async fn set_disk_id(&self, id: Option<Uuid>) -> Result<()> {
let mut lock = self.id.lock().await; let mut lock = self.id.lock().await;
@@ -115,7 +142,7 @@ impl DiskAPI for RemoteDisk {
Ok(()) Ok(())
} }
async fn read_all(&self, volume: &str, path: &str) -> Result<Bytes> { async fn read_all(&self, volume: &str, path: &str) -> Result<Vec<u8>> {
info!("read_all"); info!("read_all");
let mut client = self.get_client_v2().await?; let mut client = self.get_client_v2().await?;
let request = Request::new(ReadAllRequest { let request = Request::new(ReadAllRequest {
@@ -132,7 +159,7 @@ impl DiskAPI for RemoteDisk {
return Err(DiskError::FileNotFound.into()); return Err(DiskError::FileNotFound.into());
} }
Ok(Bytes::from(response.data)) Ok(response.data)
} }
async fn write_all(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()> { async fn write_all(&self, volume: &str, path: &str, data: Vec<u8>) -> Result<()> {
@@ -173,7 +200,26 @@ impl DiskAPI for RemoteDisk {
Ok(()) Ok(())
} }
async fn rename_part(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str, meta: Vec<u8>) -> Result<()> {
info!("rename_part");
let mut client = self.get_client_v2().await?;
let request = Request::new(RenamePartRequst {
disk: self.root.to_string_lossy().to_string(),
src_volume: src_volume.to_string(),
src_path: src_path.to_string(),
dst_volume: dst_volume.to_string(),
dst_path: dst_path.to_string(),
meta,
});
let response = client.rename_part(request).await?.into_inner();
if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
}
Ok(())
}
async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> { async fn rename_file(&self, src_volume: &str, src_path: &str, dst_volume: &str, dst_path: &str) -> Result<()> {
info!("rename_file"); info!("rename_file");
let mut client = self.get_client_v2().await?; let mut client = self.get_client_v2().await?;
@@ -253,11 +299,11 @@ impl DiskAPI for RemoteDisk {
}); });
let response = client.walk_dir(request).await?.into_inner(); let response = client.walk_dir(request).await?.into_inner();
if !response.success { if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string()))); return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
} }
let entries = response let entries = response
.meta_cache_entry .meta_cache_entry
.into_iter() .into_iter()
@@ -288,7 +334,7 @@ impl DiskAPI for RemoteDisk {
}); });
let response = client.rename_data(request).await?.into_inner(); let response = client.rename_data(request).await?.into_inner();
if !response.success { if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string()))); return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
} }
@@ -363,7 +409,7 @@ impl DiskAPI for RemoteDisk {
}); });
let response = client.stat_volume(request).await?.into_inner(); let response = client.stat_volume(request).await?.into_inner();
if !response.success { if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string()))); return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
} }
@@ -373,6 +419,47 @@ impl DiskAPI for RemoteDisk {
Ok(volume_info) Ok(volume_info)
} }
async fn delete_paths(&self, volume: &str, paths: &[&str]) -> Result<()> {
info!("delete_paths");
let paths = paths.iter().map(|s| s.to_string()).collect::<Vec<String>>();
let mut client = self.get_client_v2().await?;
let request = Request::new(DeletePathsRequest {
disk: self.root.to_string_lossy().to_string(),
volume: volume.to_string(),
paths,
});
let response = client.delete_paths(request).await?.into_inner();
if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
}
Ok(())
}
async fn update_metadata(&self, volume: &str, path: &str, fi: FileInfo, opts: UpdateMetadataOpts) -> Result<()> {
info!("update_metadata");
let file_info = serde_json::to_string(&fi)?;
let opts = serde_json::to_string(&opts)?;
let mut client = self.get_client_v2().await?;
let request = Request::new(UpdateMetadataRequest {
disk: self.root.to_string_lossy().to_string(),
volume: volume.to_string(),
path: path.to_string(),
file_info,
opts,
});
let response = client.update_metadata(request).await?.into_inner();
if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
}
Ok(())
}
async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> { async fn write_metadata(&self, _org_volume: &str, volume: &str, path: &str, fi: FileInfo) -> Result<()> {
info!("write_metadata"); info!("write_metadata");
let file_info = serde_json::to_string(&fi)?; let file_info = serde_json::to_string(&fi)?;
@@ -443,7 +530,38 @@ impl DiskAPI for RemoteDisk {
Ok(raw_file_info) Ok(raw_file_info)
} }
async fn delete_version(
&self,
volume: &str,
path: &str,
fi: FileInfo,
force_del_marker: bool,
opts: DeleteOptions,
) -> Result<RawFileInfo> {
info!("delete_version");
let file_info = serde_json::to_string(&fi)?;
let opts = serde_json::to_string(&opts)?;
let mut client = self.get_client_v2().await?;
let request = Request::new(DeleteVersionRequest {
disk: self.root.to_string_lossy().to_string(),
volume: volume.to_string(),
path: path.to_string(),
file_info,
force_del_marker,
opts,
});
let response = client.delete_version(request).await?.into_inner();
if !response.success {
return Err(Error::from_string(response.error_info.unwrap_or("".to_string())));
}
let raw_file_info = serde_json::from_str::<RawFileInfo>(&response.raw_file_info)?;
Ok(raw_file_info)
}
async fn delete_versions( async fn delete_versions(
&self, &self,
volume: &str, volume: &str,
+30 -24
View File
@@ -1,4 +1,5 @@
use crate::error::{Error, Result, StdError}; use crate::error::{Error, Result, StdError};
use crate::quorum::{object_op_ignored_errs, reduce_write_quorum_errs};
use bytes::Bytes; use bytes::Bytes;
use futures::future::join_all; use futures::future::join_all;
use futures::{Stream, StreamExt}; use futures::{Stream, StreamExt};
@@ -25,7 +26,7 @@ pub struct Erasure {
impl Erasure { impl Erasure {
pub fn new(data_shards: usize, parity_shards: usize, block_size: usize) -> Self { pub fn new(data_shards: usize, parity_shards: usize, block_size: usize) -> Self {
warn!( debug!(
"Erasure new data_shards {},parity_shards {} block_size {} ", "Erasure new data_shards {},parity_shards {} block_size {} ",
data_shards, parity_shards, block_size data_shards, parity_shards, block_size
); );
@@ -46,17 +47,17 @@ impl Erasure {
pub async fn encode<S>( pub async fn encode<S>(
&self, &self,
body: S, body: S,
writers: &mut [FileWriter], writers: &mut [Option<FileWriter>],
// block_size: usize, // block_size: usize,
total_size: usize, total_size: usize,
_write_quorum: usize, write_quorum: usize,
) -> Result<usize> ) -> Result<usize>
where where
S: Stream<Item = Result<Bytes, StdError>> + Send + Sync + 'static, S: Stream<Item = Result<Bytes, StdError>> + Send + Sync + 'static,
{ {
let mut stream = ChunkedStream::new(body, total_size, self.block_size, false); let mut stream = ChunkedStream::new(body, total_size, self.block_size, false);
let mut total: usize = 0; let mut total: usize = 0;
let mut idx = 0; // let mut idx = 0;
while let Some(result) = stream.next().await { while let Some(result) = stream.next().await {
match result { match result {
Ok(data) => { Ok(data) => {
@@ -67,36 +68,41 @@ impl Erasure {
break; break;
} }
idx += 1; // idx += 1;
debug!("encode {} get data {}", idx, data.len()); // debug!("encode {} get data {}", idx, data.len());
let blocks = self.encode_data(data.as_ref())?; let blocks = self.encode_data(data.as_ref())?;
debug!( // debug!(
"encode shard {} size: {}/{} from block_size {}, total_size {} ", // "encode shard {} size: {}/{} from block_size {}, total_size {} ",
idx, // idx,
blocks[0].len(), // blocks[0].len(),
blocks.len(), // blocks.len(),
data.len(), // data.len(),
total_size // total_size
); // );
let mut errs = Vec::new(); let mut errs = Vec::new();
for (i, w) in writers.iter_mut().enumerate() { for (i, w) in writers.iter_mut().enumerate() {
match w.write(blocks[i].as_ref()).await { if w.is_none() {
continue;
}
match w.as_mut().unwrap().write(blocks[i].as_ref()).await {
Ok(_) => errs.push(None), Ok(_) => errs.push(None),
Err(e) => errs.push(Some(e)), Err(e) => errs.push(Some(e)),
} }
} }
// debug!("{} encode_data write errs:{:?}", self.id, errs); let none_count = errs.iter().filter(|&x| x.is_none()).count();
// // TODO: reduceWriteQuorumErrs if none_count >= write_quorum {
// for err in errs.iter() { continue;
// if err.is_some() { }
// return Err(Error::msg("message"));
// } if let Some(err) = reduce_write_quorum_errs(&errs, object_op_ignored_errs().as_ref(), write_quorum) {
// } warn!("Erasure encode errs {:?}", &errs);
return Err(err);
}
} }
Err(e) => return Err(Error::from_std_error(e)), Err(e) => return Err(Error::from_std_error(e)),
} }
@@ -154,7 +160,7 @@ impl Erasure {
break; break;
} }
debug!("decode {} block_offset {},block_length {} ", block_idx, block_offset, block_length); // debug!("decode {} block_offset {},block_length {} ", block_idx, block_offset, block_length);
let mut bufs = reader.read().await?; let mut bufs = reader.read().await?;
@@ -168,7 +174,7 @@ impl Erasure {
bytes_writed += writed_n; bytes_writed += writed_n;
debug!("decode {} writed_n {}, total_writed: {} ", block_idx, writed_n, bytes_writed); // debug!("decode {} writed_n {}, total_writed: {} ", block_idx, writed_n, bytes_writed);
} }
if bytes_writed != length { if bytes_writed != length {
+25
View File
@@ -1,5 +1,9 @@
use std::io;
use tracing_error::{SpanTrace, SpanTraceStatus}; use tracing_error::{SpanTrace, SpanTraceStatus};
use crate::disk::error::{clone_disk_err, DiskError};
pub type StdError = Box<dyn std::error::Error + Send + Sync + 'static>; pub type StdError = Box<dyn std::error::Error + Send + Sync + 'static>;
pub type Result<T = (), E = Error> = std::result::Result<T, E>; pub type Result<T = (), E = Error> = std::result::Result<T, E>;
@@ -61,6 +65,14 @@ impl Error {
pub fn downcast_mut<T: std::error::Error + 'static>(&mut self) -> Option<&mut T> { pub fn downcast_mut<T: std::error::Error + 'static>(&mut self) -> Option<&mut T> {
self.inner.downcast_mut() self.inner.downcast_mut()
} }
pub fn to_io_err(&self) -> Option<io::Error> {
if let Some(e) = self.downcast_ref::<io::Error>() {
Some(io::Error::new(e.kind(), e.to_string()))
} else {
None
}
}
} }
impl<T: std::error::Error + Send + Sync + 'static> From<T> for Error { impl<T: std::error::Error + Send + Sync + 'static> From<T> for Error {
@@ -80,3 +92,16 @@ impl std::fmt::Display for Error {
Ok(()) Ok(())
} }
} }
impl Clone for Error {
fn clone(&self) -> Self {
if let Some(e) = self.downcast_ref::<DiskError>() {
clone_disk_err(e)
} else if let Some(e) = self.downcast_ref::<io::Error>() {
Error::new(io::Error::new(e.kind(), e.to_string()))
} else {
// TODO: 优化其他类型
Error::msg(self.to_string())
}
}
}
+2
View File
@@ -7,8 +7,10 @@ pub mod erasure;
pub mod error; pub mod error;
mod file_meta; mod file_meta;
pub mod peer; pub mod peer;
mod quorum;
pub mod set_disk; pub mod set_disk;
mod sets; mod sets;
mod storage_class;
pub mod store; pub mod store;
pub mod store_api; pub mod store_api;
mod store_init; mod store_init;
+1
View File
@@ -262,6 +262,7 @@ impl PeerS3Client for LocalPeerS3Client {
} }
} }
warn!("list_bucket ress {:?}", &ress);
warn!("list_bucket errs {:?}", &errs); warn!("list_bucket errs {:?}", &errs);
let mut uniq_map: HashMap<&String, &VolumeInfo> = HashMap::new(); let mut uniq_map: HashMap<&String, &VolumeInfo> = HashMap::new();
+127
View File
@@ -0,0 +1,127 @@
use crate::{disk::error::DiskError, error::Error};
use std::{collections::HashMap, fmt::Debug};
// pub type CheckErrorFn = fn(e: &Error) -> bool;
pub trait CheckErrorFn: Debug + Send + Sync + 'static {
fn is(&self, e: &Error) -> bool;
}
#[derive(Debug, thiserror::Error)]
pub enum QuorumError {
#[error("Read quorum not met")]
Read,
#[error("disk not found")]
Write,
}
pub fn base_ignored_errs() -> Vec<Box<dyn CheckErrorFn>> {
vec![
Box::new(DiskError::DiskNotFound),
Box::new(DiskError::FaultyDisk),
Box::new(DiskError::FaultyRemoteDisk),
]
}
// object_op_ignored_errs
pub fn object_op_ignored_errs() -> Vec<Box<dyn CheckErrorFn>> {
let mut base = base_ignored_errs();
let ext: Vec<Box<dyn CheckErrorFn>> = vec![
// Box::new(DiskError::DiskNotFound),
// Box::new(DiskError::FaultyDisk),
// Box::new(DiskError::FaultyRemoteDisk),
Box::new(DiskError::DiskAccessDenied),
Box::new(DiskError::UnformattedDisk),
Box::new(DiskError::DiskOngoingReq),
];
base.extend(ext);
base
}
// 用于检查错误是否被忽略的函数
fn is_err_ignored(err: &Error, ignored_errs: &Vec<Box<dyn CheckErrorFn>>) -> bool {
ignored_errs.iter().any(|ignored_err| ignored_err.is(err))
}
// 减少错误数量并返回出现次数最多的错误
fn reduce_errs(errs: &Vec<Option<Error>>, ignored_errs: &Vec<Box<dyn CheckErrorFn>>) -> (usize, Option<Error>) {
let mut error_counts: HashMap<String, usize> = HashMap::new();
let mut error_map: HashMap<String, usize> = HashMap::new(); // 存err位置
let nil = "nil".to_string();
for (i, operr) in errs.iter().enumerate() {
if operr.is_none() {
*error_counts.entry(nil.clone()).or_insert(0) += 1;
let _ = *error_map.entry(nil.clone()).or_insert(i);
continue;
}
let err = operr.as_ref().unwrap();
if is_err_ignored(err, &ignored_errs) {
continue;
}
let errstr = err.to_string();
let _ = *error_map.entry(errstr.clone()).or_insert(i);
*error_counts.entry(errstr.clone()).or_insert(0) += 1;
}
let mut max = 0;
let mut max_err = nil.clone();
for (&ref err, &count) in error_counts.iter() {
if count > max || (count == max && *err == nil) {
max = count;
max_err = err.clone();
}
}
if let Some(&c) = error_counts.get(&max_err) {
if let Some(&err_idx) = error_map.get(&max_err) {
let err = errs[err_idx].clone();
return (c, err);
}
return (c, None);
}
(0, None)
}
// 根据quorum验证错误数量
fn reduce_quorum_errs(
errs: &Vec<Option<Error>>,
ignored_errs: &Vec<Box<dyn CheckErrorFn>>,
quorum: usize,
quorum_err: QuorumError,
) -> Option<Error> {
let (max_count, max_err) = reduce_errs(errs, ignored_errs);
if max_count >= quorum {
max_err
} else {
Some(Error::new(quorum_err))
}
}
// 根据读quorum验证错误数量
// 返回最大错误数量的下标,或QuorumError
pub fn reduce_read_quorum_errs(
errs: &Vec<Option<Error>>,
ignored_errs: &Vec<Box<dyn CheckErrorFn>>,
read_quorum: usize,
) -> Option<Error> {
reduce_quorum_errs(errs, ignored_errs, read_quorum, QuorumError::Read)
}
// 根据写quorum验证错误数量
// 返回最大错误数量的下标,或QuorumError
pub fn reduce_write_quorum_errs(
errs: &Vec<Option<Error>>,
ignored_errs: &Vec<Box<dyn CheckErrorFn>>,
write_quorum: usize,
) -> Option<Error> {
reduce_quorum_errs(errs, ignored_errs, write_quorum, QuorumError::Write)
}
+72 -18
View File
@@ -1,10 +1,5 @@
#![allow(clippy::map_entry)] #![allow(clippy::map_entry)]
use std::collections::HashMap; use std::{collections::HashMap, sync::Arc};
use futures::future::join_all;
use http::HeaderMap;
use tracing::warn;
use uuid::Uuid;
use crate::{ use crate::{
disk::{ disk::{
@@ -22,13 +17,21 @@ use crate::{
}, },
utils::hash, utils::hash,
}; };
use futures::future::join_all;
use http::HeaderMap;
use tokio::sync::RwLock;
use tokio::time::Duration;
use tokio_util::sync::CancellationToken;
use tracing::info;
use tracing::warn;
use uuid::Uuid;
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct Sets { pub struct Sets {
pub id: Uuid, pub id: Uuid,
// pub sets: Vec<Objects>, // pub sets: Vec<Objects>,
// pub disk_set: Vec<Vec<Option<DiskStore>>>, // [set_count_idx][set_drive_count_idx] = disk_idx // pub disk_set: Vec<Vec<Option<DiskStore>>>, // [set_count_idx][set_drive_count_idx] = disk_idx
pub disk_set: Vec<SetDisks>, // [set_count_idx][set_drive_count_idx] = disk_idx pub disk_set: Vec<Arc<SetDisks>>, // [set_count_idx][set_drive_count_idx] = disk_idx
pub pool_idx: usize, pub pool_idx: usize,
pub endpoints: PoolEndpoints, pub endpoints: PoolEndpoints,
pub format: FormatV3, pub format: FormatV3,
@@ -36,6 +39,7 @@ pub struct Sets {
pub set_count: usize, pub set_count: usize,
pub set_drive_count: usize, pub set_drive_count: usize,
pub distribution_algo: DistributionAlgoVersion, pub distribution_algo: DistributionAlgoVersion,
ctx: CancellationToken,
} }
impl Sets { impl Sets {
@@ -45,7 +49,7 @@ impl Sets {
fm: &FormatV3, fm: &FormatV3,
pool_idx: usize, pool_idx: usize,
partiy_count: usize, partiy_count: usize,
) -> Result<Self> { ) -> Result<Arc<Self>> {
let set_count = fm.erasure.sets.len(); let set_count = fm.erasure.sets.len();
let set_drive_count = fm.erasure.sets[0].len(); let set_drive_count = fm.erasure.sets[0].len();
@@ -53,9 +57,14 @@ impl Sets {
for i in 0..set_count { for i in 0..set_count {
let mut set_drive = Vec::with_capacity(set_drive_count); let mut set_drive = Vec::with_capacity(set_drive_count);
let mut set_endpoints = Vec::with_capacity(set_drive_count);
for j in 0..set_drive_count { for j in 0..set_drive_count {
let idx = i * set_drive_count + j; let idx = i * set_drive_count + j;
let mut disk = disks[idx].clone(); let mut disk = disks[idx].clone();
let endpoint = endpoints.endpoints.as_ref().get(idx).cloned();
set_endpoints.push(endpoint);
if disk.is_none() { if disk.is_none() {
warn!("sets new set_drive {}-{} is none", i, j); warn!("sets new set_drive {}-{} is none", i, j);
set_drive.push(None); set_drive.push(None);
@@ -79,7 +88,7 @@ impl Sets {
disk = local_disk; disk = local_disk;
} }
if let Some(_disk_id) = disk.as_ref().unwrap().get_disk_id().await { if let Some(_disk_id) = disk.as_ref().unwrap().get_disk_id().await? {
set_drive.push(disk); set_drive.push(disk);
} else { } else {
warn!("sets new set_drive {}-{} get_disk_id is none", i, j); warn!("sets new set_drive {}-{} get_disk_id is none", i, j);
@@ -87,20 +96,22 @@ impl Sets {
} }
} }
warn!("sets new set_drive {:?}", &set_drive); // warn!("sets new set_drive {:?}", &set_drive);
let set_disks = SetDisks { let set_disks = SetDisks {
disks: set_drive, disks: RwLock::new(set_drive),
set_drive_count, set_drive_count,
parity_count: partiy_count, default_parity_count: partiy_count,
set_index: i, set_index: i,
pool_index: pool_idx, pool_index: pool_idx,
set_endpoints,
format: fm.clone(),
}; };
disk_set.push(set_disks); disk_set.push(Arc::new(set_disks));
} }
let sets = Self { let sets = Arc::new(Self {
id: fm.id, id: fm.id,
// sets: todo!(), // sets: todo!(),
disk_set, disk_set,
@@ -111,15 +122,58 @@ impl Sets {
set_count, set_count,
set_drive_count, set_drive_count,
distribution_algo: fm.erasure.distribution_algo.clone(), distribution_algo: fm.erasure.distribution_algo.clone(),
}; ctx: CancellationToken::new(),
});
let asets = sets.clone();
tokio::spawn(async move { asets.monitor_and_connect_endpoints().await });
Ok(sets) Ok(sets)
} }
pub fn get_disks(&self, set_idx: usize) -> SetDisks {
pub async fn monitor_and_connect_endpoints(&self) {
tokio::time::sleep(Duration::from_secs(5)).await;
info!("start monitor_and_connect_endpoints");
self.connect_disks().await;
// TODO: config interval
let mut interval = tokio::time::interval(tokio::time::Duration::from_secs(15 * 3));
let cloned_token = self.ctx.clone();
loop {
tokio::select! {
_= interval.tick()=>{
// debug!("tick...");
self.connect_disks().await;
interval.reset();
},
_ = cloned_token.cancelled() => {
warn!("ctx cancelled");
break;
}
}
}
warn!("monitor_and_connect_endpoints exit");
}
async fn connect_disks(&self) {
// debug!("start connect_disks ...");
for set in self.disk_set.iter() {
set.connect_disks().await;
}
// debug!("done connect_disks ...");
}
pub fn get_disks(&self, set_idx: usize) -> Arc<SetDisks> {
self.disk_set[set_idx].clone() self.disk_set[set_idx].clone()
} }
pub fn get_disks_by_key(&self, key: &str) -> SetDisks { pub fn get_disks_by_key(&self, key: &str) -> Arc<SetDisks> {
self.get_disks(self.get_hashed_set_index(key)) self.get_disks(self.get_hashed_set_index(key))
} }
@@ -314,7 +368,7 @@ impl StorageAPI for Sets {
.get_object_reader(bucket, object, range, h, opts) .get_object_reader(bucket, object, range, h, opts)
.await .await
} }
async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<()> { async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
self.get_disks_by_key(object).put_object(bucket, object, data, opts).await self.get_disks_by_key(object).put_object(bucket, object, data, opts).await
} }
+28
View File
@@ -0,0 +1,28 @@
// use crate::error::{Error, Result};
// default_partiy_count 默认配置,根据磁盘总数分配校验磁盘数量
pub fn default_partiy_count(drive: usize) -> usize {
match drive {
1 => 0,
2 | 3 => 1,
4 | 5 => 2,
6 | 7 => 3,
_ => 4,
}
}
// Define the minimum number of parity drives required.
// const MIN_PARITY_DRIVES: usize = 0;
// // ValidateParity validates standard storage class parity.
// pub fn validate_parity(ss_parity: usize, set_drive_count: usize) -> Result<()> {
// // if ss_parity > 0 && ss_parity < MIN_PARITY_DRIVES {
// // return Err(Error::msg(format!("parity {} 应该大于等于 {}", ss_parity, MIN_PARITY_DRIVES)));
// // }
// if ss_parity > set_drive_count / 2 {
// return Err(Error::msg(format!("parity {} 应该小于等于 {}", ss_parity, set_drive_count / 2)));
// }
// Ok(())
// }
+60 -68
View File
@@ -1,13 +1,12 @@
#![allow(clippy::map_entry)] #![allow(clippy::map_entry)]
use crate::{ use crate::{
bucket_meta::BucketMetadata, bucket_meta::BucketMetadata,
disk::{ disk::{error::DiskError, new_disk, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET},
error::DiskError, new_disk, DeleteOptions, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET,
},
endpoints::{EndpointServerPools, SetupType}, endpoints::{EndpointServerPools, SetupType},
error::{Error, Result}, error::{Error, Result},
peer::S3PeerSys, peer::S3PeerSys,
sets::Sets, sets::Sets,
storage_class::default_partiy_count,
store_api::{ store_api::{
BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec, BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec,
ListObjectsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, ObjectToDelete, ListObjectsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, ObjectToDelete,
@@ -25,11 +24,9 @@ use std::{
time::Duration, time::Duration,
}; };
use time::OffsetDateTime; use time::OffsetDateTime;
use tokio::{ use tokio::sync::Semaphore;
fs, use tokio::{fs, sync::RwLock};
sync::{RwLock, Semaphore}, use tracing::{debug, info};
};
use tracing::{debug, info, warn};
use uuid::Uuid; use uuid::Uuid;
use lazy_static::lazy_static; use lazy_static::lazy_static;
@@ -155,7 +152,7 @@ pub struct ECStore {
pub id: uuid::Uuid, pub id: uuid::Uuid,
// pub disks: Vec<DiskStore>, // pub disks: Vec<DiskStore>,
pub disk_map: HashMap<usize, Vec<Option<DiskStore>>>, pub disk_map: HashMap<usize, Vec<Option<DiskStore>>>,
pub pools: Vec<Sets>, pub pools: Vec<Arc<Sets>>,
pub peer_sys: S3PeerSys, pub peer_sys: S3PeerSys,
// pub local_disks: Vec<DiskStore>, // pub local_disks: Vec<DiskStore>,
} }
@@ -176,11 +173,13 @@ impl ECStore {
let mut local_disks = Vec::new(); let mut local_disks = Vec::new();
info!("endpoint_pools: {:?}", endpoint_pools); debug!("endpoint_pools: {:?}", endpoint_pools);
for (i, pool_eps) in endpoint_pools.as_ref().iter().enumerate() { for (i, pool_eps) in endpoint_pools.as_ref().iter().enumerate() {
// TODO: read from config parseStorageClass // TODO: read from config parseStorageClass
let partiy_count = store_init::default_partiy_count(pool_eps.drives_per_set); let partiy_count = default_partiy_count(pool_eps.drives_per_set);
// validate_parity(partiy_count, pool_eps.drives_per_set)?;
let (disks, errs) = crate::store_init::init_disks( let (disks, errs) = crate::store_init::init_disks(
&pool_eps.endpoints, &pool_eps.endpoints,
@@ -229,7 +228,6 @@ impl ECStore {
} }
let sets = Sets::new(disks.clone(), pool_eps, &fm, i, partiy_count).await?; let sets = Sets::new(disks.clone(), pool_eps, &fm, i, partiy_count).await?;
pools.push(sets); pools.push(sets);
disk_map.insert(i, disks); disk_map.insert(i, disks);
@@ -291,62 +289,54 @@ impl ECStore {
for sets in self.pools.iter() { for sets in self.pools.iter() {
for set in sets.disk_set.iter() { for set in sets.disk_set.iter() {
for disk in set.disks.iter() { futures.push(set.walk_dir(&opts));
if disk.is_none() {
continue;
}
let disk = disk.as_ref().unwrap();
let opts = opts.clone();
// let mut wr = &mut wr;
futures.push(disk.walk_dir(opts));
// tokio::spawn(async move { disk.walk_dir(opts, wr).await });
}
} }
} }
let results = join_all(futures).await; let results = join_all(futures).await;
let mut errs = Vec::new(); // let mut errs = Vec::new();
let mut ress = Vec::new(); let mut ress = Vec::new();
let mut uniq = HashSet::new(); let mut uniq = HashSet::new();
for res in results { for (disks_ress, _disks_errs) in results {
match res { for (_i, disks_res) in disks_ress.iter().enumerate() {
Ok(entrys) => { if disks_res.is_none() {
for entry in entrys { // TODO handle errs
if !uniq.contains(&entry.name) { continue;
uniq.insert(entry.name.clone()); }
// TODO: 过滤 let entrys = disks_res.as_ref().unwrap();
if opts.limit > 0 && ress.len() as i32 >= opts.limit {
return Ok(ress);
}
if entry.is_object() { for entry in entrys {
let fi = entry.to_fileinfo(&opts.bucket)?; if !uniq.contains(&entry.name) {
if let Some(f) = fi { uniq.insert(entry.name.clone());
ress.push(f.into_object_info(&opts.bucket, &entry.name, false)); // TODO: 过滤
} if opts.limit > 0 && ress.len() as i32 >= opts.limit {
continue; return Ok(ress);
} }
if entry.is_dir() { if entry.is_object() {
ress.push(ObjectInfo { let fi = entry.to_fileinfo(&opts.bucket)?;
is_dir: true, if let Some(f) = fi {
bucket: opts.bucket.clone(), ress.push(f.to_object_info(&opts.bucket, &entry.name, false));
name: entry.name,
..Default::default()
});
} }
continue;
}
if entry.is_dir() {
ress.push(ObjectInfo {
is_dir: true,
bucket: opts.bucket.clone(),
name: entry.name.clone(),
..Default::default()
});
} }
} }
errs.push(None);
} }
Err(e) => errs.push(Some(e)),
} }
} }
warn!("list_merged errs {:?}", errs); // warn!("list_merged errs {:?}", errs);
Ok(ress) Ok(ress)
} }
@@ -355,21 +345,23 @@ impl ECStore {
let mut futures = Vec::new(); let mut futures = Vec::new();
for sets in self.pools.iter() { for sets in self.pools.iter() {
for set in sets.disk_set.iter() { for set in sets.disk_set.iter() {
for disk in set.disks.iter() { futures.push(set.delete_all(bucket, prefix));
if disk.is_none() { // let disks = set.disks.read().await;
continue; // let dd = disks.clone();
} // for disk in dd {
// if disk.is_none() {
let disk = disk.as_ref().unwrap(); // continue;
futures.push(disk.delete( // }
bucket, // // let disk = disk.as_ref().unwrap().clone();
prefix, // // futures.push(disk.delete(
DeleteOptions { // // bucket,
recursive: true, // // prefix,
immediate: false, // // DeleteOptions {
}, // // recursive: true,
)); // // immediate: false,
} // // },
// // ));
// }
} }
} }
let results = join_all(futures).await; let results = join_all(futures).await;
@@ -402,7 +394,7 @@ impl ECStore {
} }
async fn internal_get_pool_info_existing_with_opts( async fn internal_get_pool_info_existing_with_opts(
pools: &[Sets], pools: &[Arc<Sets>],
bucket: &str, bucket: &str,
object: &str, object: &str,
opts: &ObjectOptions, opts: &ObjectOptions,
@@ -755,7 +747,7 @@ impl StorageAPI for ECStore {
unimplemented!() unimplemented!()
} }
async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<()> { async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo> {
// checkPutObjectArgs // checkPutObjectArgs
let object = utils::path::encode_dir_object(object); let object = utils::path::encode_dir_object(object);
+13 -3
View File
@@ -27,8 +27,10 @@ pub struct FileInfo {
pub fresh: bool, // indicates this is a first time call to write FileInfo. pub fresh: bool, // indicates this is a first time call to write FileInfo.
pub parts: Vec<ObjectPartInfo>, pub parts: Vec<ObjectPartInfo>,
pub is_latest: bool, pub is_latest: bool,
#[serde(skip_serializing_if = "Option::is_none", default)] // #[serde(skip_serializing_if = "Option::is_none", default)]
pub tags: Option<HashMap<String, String>>, pub tags: Option<HashMap<String, String>>,
pub metadata: Option<HashMap<String, String>>,
pub num_versions: usize,
} }
// impl Default for FileInfo { // impl Default for FileInfo {
@@ -96,6 +98,14 @@ impl FileInfo {
false false
} }
pub fn get_etag(&self) -> Option<String> {
if let Some(meta) = &self.metadata {
meta.get("etag").cloned()
} else {
None
}
}
pub fn write_quorum(&self, quorum: usize) -> usize { pub fn write_quorum(&self, quorum: usize) -> usize {
if self.deleted { if self.deleted {
return quorum; return quorum;
@@ -141,7 +151,7 @@ impl FileInfo {
self.parts.sort_by(|a, b| a.number.cmp(&b.number)); self.parts.sort_by(|a, b| a.number.cmp(&b.number));
} }
pub fn into_object_info(&self, bucket: &str, object: &str, _versioned: bool) -> ObjectInfo { pub fn to_object_info(&self, bucket: &str, object: &str, _versioned: bool) -> ObjectInfo {
ObjectInfo { ObjectInfo {
bucket: bucket.to_string(), bucket: bucket.to_string(),
name: object.to_string(), name: object.to_string(),
@@ -533,7 +543,7 @@ pub trait StorageAPI {
h: HeaderMap, h: HeaderMap,
opts: &ObjectOptions, opts: &ObjectOptions,
) -> Result<GetObjectReader>; ) -> Result<GetObjectReader>;
async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<()>; async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<ObjectInfo>;
async fn put_object_part( async fn put_object_part(
&self, &self,
bucket: &str, bucket: &str,
+18 -35
View File
@@ -13,7 +13,7 @@ use std::{
fmt::Debug, fmt::Debug,
}; };
use tracing::warn; use tracing::{debug, warn};
use uuid::Uuid; use uuid::Uuid;
pub async fn init_disks(eps: &Endpoints, opt: &DiskOption) -> (Vec<Option<DiskStore>>, Vec<Option<Error>>) { pub async fn init_disks(eps: &Endpoints, opt: &DiskOption) -> (Vec<Option<DiskStore>>, Vec<Option<Error>>) {
@@ -50,13 +50,11 @@ pub async fn connect_load_init_formats(
set_drive_count: usize, set_drive_count: usize,
deployment_id: Option<Uuid>, deployment_id: Option<Uuid>,
) -> Result<FormatV3, Error> { ) -> Result<FormatV3, Error> {
warn!("connect_load_init_formats id: {:?}, first_disk: {}", deployment_id, first_disk);
let (formats, errs) = load_format_erasure_all(disks, false).await; let (formats, errs) = load_format_erasure_all(disks, false).await;
DiskError::check_disk_fatal_errs(&errs)?; debug!("load_format_erasure_all errs {:?}", &errs);
warn!("load_format_erasure_all errs {:?}", &errs); DiskError::check_disk_fatal_errs(&errs)?;
check_format_erasure_values(&formats, set_drive_count)?; check_format_erasure_values(&formats, set_drive_count)?;
@@ -65,8 +63,9 @@ pub async fn connect_load_init_formats(
// new format and save // new format and save
let fms = init_format_erasure(disks, set_count, set_drive_count, deployment_id); let fms = init_format_erasure(disks, set_count, set_drive_count, deployment_id);
let _errs = save_format_file_all(disks, &fms).await; let errs = save_format_file_all(disks, &fms).await;
debug!("save_format_file_all errs {:?}", &errs);
// TODO: check quorum // TODO: check quorum
// reduceWriteQuorumErrs(&errs)?; // reduceWriteQuorumErrs(&errs)?;
@@ -131,15 +130,8 @@ fn get_format_erasure_in_quorum(formats: &[Option<FormatV3>]) -> Result<FormatV3
let (max_drives, max_count) = countmap.iter().max_by_key(|&(_, c)| c).unwrap_or((&0, &0)); let (max_drives, max_count) = countmap.iter().max_by_key(|&(_, c)| c).unwrap_or((&0, &0));
warn!("get_format_erasure_in_quorum fi: {:?}", &formats);
if *max_drives == 0 || *max_count <= formats.len() / 2 { if *max_drives == 0 || *max_count <= formats.len() / 2 {
warn!( warn!("get_format_erasure_in_quorum fi: {:?}", &formats);
"*max_drives == 0 || *max_count < formats.len() / 2, {} || {}<{}",
max_drives,
max_count,
formats.len() / 2
);
return Err(Error::new(ErasureError::ErasureReadQuorum)); return Err(Error::new(ErasureError::ErasureReadQuorum));
} }
@@ -189,26 +181,22 @@ fn check_format_erasure_value(format: &FormatV3) -> Result<()> {
Ok(()) Ok(())
} }
pub fn default_partiy_count(drive: usize) -> usize {
match drive {
1 => 0,
2 | 3 => 1,
4 | 5 => 2,
6 | 7 => 3,
_ => 4,
}
}
// load_format_erasure_all 读取所有foramt.json // load_format_erasure_all 读取所有foramt.json
async fn load_format_erasure_all(disks: &[Option<DiskStore>], heal: bool) -> (Vec<Option<FormatV3>>, Vec<Option<Error>>) { async fn load_format_erasure_all(disks: &[Option<DiskStore>], heal: bool) -> (Vec<Option<FormatV3>>, Vec<Option<Error>>) {
let mut futures = Vec::with_capacity(disks.len()); let mut futures = Vec::with_capacity(disks.len());
for disk in disks.iter() {
futures.push(read_format_file(disk, heal));
}
let mut datas = Vec::with_capacity(disks.len()); let mut datas = Vec::with_capacity(disks.len());
let mut errors = Vec::with_capacity(disks.len()); let mut errors = Vec::with_capacity(disks.len());
for disk in disks.iter() {
if disk.is_none() {
datas.push(None);
errors.push(Some(Error::new(DiskError::DiskNotFound)));
}
let disk = disk.as_ref().unwrap();
futures.push(load_format_erasure(disk, heal));
}
let results = join_all(futures).await; let results = join_all(futures).await;
let mut i = 0; let mut i = 0;
for result in results { for result in results {
@@ -233,12 +221,7 @@ async fn load_format_erasure_all(disks: &[Option<DiskStore>], heal: bool) -> (Ve
(datas, errors) (datas, errors)
} }
async fn read_format_file(disk: &Option<DiskStore>, _heal: bool) -> Result<FormatV3, Error> { pub async fn load_format_erasure(disk: &DiskStore, _heal: bool) -> Result<FormatV3, Error> {
if disk.is_none() {
return Err(Error::new(DiskError::DiskNotFound));
}
let disk = disk.as_ref().unwrap();
let data = disk let data = disk
.read_all(RUSTFS_META_BUCKET, FORMAT_CONFIG_FILE) .read_all(RUSTFS_META_BUCKET, FORMAT_CONFIG_FILE)
.await .await
@@ -249,7 +232,7 @@ async fn read_format_file(disk: &Option<DiskStore>, _heal: bool) -> Result<Forma
None => e, None => e,
})?; })?;
let fm = FormatV3::try_from(data.as_ref())?; let fm = FormatV3::try_from(data.as_slice())?;
// TODO: heal // TODO: heal
+132
View File
@@ -0,0 +1,132 @@
use std::{fs::Metadata, path::Path};
use tokio::{
fs::{self, File},
io,
};
#[cfg(not(windows))]
pub fn same_file(f1: &Metadata, f2: &Metadata) -> bool {
use std::os::unix::fs::MetadataExt;
if f1.dev() != f2.dev() {
return false;
}
if f1.ino() != f2.ino() {
return false;
}
if f1.size() != f2.size() {
return false;
}
if f1.permissions() != f2.permissions() {
return false;
}
if f1.mtime() != f2.mtime() {
return false;
}
true
}
#[cfg(windows)]
pub fn same_file(f1: &Metadata, f2: &Metadata) -> bool {
if f1.permissions() != f2.permissions() {
return false;
}
if f1.file_type() != f2.file_type() {
return false;
}
if f1.len() != f2.len() {
return false;
}
true
}
type FileMode = usize;
pub const O_RDONLY: FileMode = 0x00000;
pub const O_WRONLY: FileMode = 0x00001;
pub const O_RDWR: FileMode = 0x00002;
pub const O_CREATE: FileMode = 0x00040;
// pub const O_EXCL: FileMode = 0x00080;
// pub const O_NOCTTY: FileMode = 0x00100;
pub const O_TRUNC: FileMode = 0x00200;
// pub const O_NONBLOCK: FileMode = 0x00800;
pub const O_APPEND: FileMode = 0x00400;
// pub const O_SYNC: FileMode = 0x01000;
// pub const O_ASYNC: FileMode = 0x02000;
// pub const O_CLOEXEC: FileMode = 0x80000;
// read: bool,
// write: bool,
// append: bool,
// truncate: bool,
// create: bool,
// create_new: bool,
pub async fn open_file(path: impl AsRef<Path>, mode: FileMode) -> io::Result<File> {
let mut opts = fs::OpenOptions::new();
match mode & (O_RDONLY | O_WRONLY | O_RDWR) {
O_RDONLY => {
opts.read(true);
}
O_WRONLY => {
opts.write(true);
}
O_RDWR => {
opts.read(true);
opts.write(true);
}
_ => (),
};
if mode & O_CREATE != 0 {
opts.create(true);
}
if mode & O_APPEND != 0 {
opts.append(true);
}
if mode & O_TRUNC != 0 {
opts.truncate(true);
}
opts.open(path.as_ref()).await
}
pub async fn access(path: impl AsRef<Path>) -> io::Result<()> {
fs::metadata(path).await?;
Ok(())
}
pub async fn lstat(path: impl AsRef<Path>) -> io::Result<Metadata> {
fs::metadata(path).await
}
pub async fn make_dir_all(path: impl AsRef<Path>) -> io::Result<()> {
fs::create_dir_all(path.as_ref()).await
}
pub async fn remove(path: impl AsRef<Path>) -> io::Result<()> {
let meta = fs::metadata(path.as_ref()).await?;
if meta.is_dir() {
fs::remove_dir(path.as_ref()).await
} else {
fs::remove_file(path.as_ref()).await
}
}
pub async fn mkdir(path: impl AsRef<Path>) -> io::Result<()> {
fs::create_dir(path.as_ref()).await
}
pub async fn rename(from: impl AsRef<Path>, to: impl AsRef<Path>) -> io::Result<()> {
fs::rename(from, to).await
}
+1
View File
@@ -1,5 +1,6 @@
pub mod crypto; pub mod crypto;
pub mod ellipses; pub mod ellipses;
pub mod fs;
pub mod hash; pub mod hash;
pub mod net; pub mod net;
pub mod path; pub mod path;
+1 -1
View File
@@ -1,6 +1,6 @@
const GLOBAL_DIR_SUFFIX: &str = "__XLDIR__"; const GLOBAL_DIR_SUFFIX: &str = "__XLDIR__";
const SLASH_SEPARATOR: &str = "/"; pub const SLASH_SEPARATOR: &str = "/";
pub fn has_suffix(s: &str, suffix: &str) -> bool { pub fn has_suffix(s: &str, suffix: &str) -> bool {
if cfg!(target_os = "windows") { if cfg!(target_os = "windows") {
+5 -32
View File
@@ -39,38 +39,11 @@ impl<'a> AsyncWrite for AppendWriter<'a> {
let mut fut = Box::pin(self.async_write(buf)); let mut fut = Box::pin(self.async_write(buf));
debug!("AsyncWrite poll_write {}, buf:{}", self.disk.id(), buf.len()); debug!("AsyncWrite poll_write {}, buf:{}", self.disk.id(), buf.len());
// while let Poll::Ready(e) = fut.as_mut().poll(cx) { let mut fut = self.get_mut().async_write(buf);
// let a = match e { match futures::future::poll_fn(|cx| fut.as_mut().poll(cx)).start(cx) {
// Ok(_) => { Ready(Ok(n)) => Ready(Ok(n)),
// debug!("Ready ok {}", self.disk.id()); Ready(Err(e)) => Ready(Err(e)),
// Poll::Ready(Ok(buf.len())) Pending => Pending,
// }
// Err(e) => {
// debug!("Ready err {}", self.disk.id());
// Poll::Ready(Err(e))
// }
// };
// return a;
// }
// Poll::Pending
match fut.as_mut().poll(cx) {
Poll::Pending => {
debug!("Pending {}", self.disk.id());
Poll::Pending
}
Poll::Ready(e) => match e {
Ok(_) => {
debug!("Ready ok {}", self.disk.id());
Poll::Ready(Ok(buf.len()))
}
Err(e) => {
debug!("Ready err {}", self.disk.id());
Poll::Ready(Err(e))
}
},
} }
} }
+157 -9
View File
@@ -1,5 +1,5 @@
use ecstore::{ use ecstore::{
disk::{DeleteOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, WalkDirOptions}, disk::{DeleteOptions, DiskStore, FileInfoVersions, ReadMultipleReq, ReadOptions, UpdateMetadataOpts, WalkDirOptions},
erasure::{ReadAt, Write}, erasure::{ReadAt, Write},
peer::{LocalPeerS3Client, PeerS3Client}, peer::{LocalPeerS3Client, PeerS3Client},
store::{all_local_disk_path, find_local_disk}, store::{all_local_disk_path, find_local_disk},
@@ -12,14 +12,16 @@ use protos::{
models::{PingBody, PingBodyBuilder}, models::{PingBody, PingBodyBuilder},
proto_gen::node_service::{ proto_gen::node_service::{
node_service_server::{NodeService as Node, NodeServiceServer as NodeServer}, node_service_server::{NodeService as Node, NodeServiceServer as NodeServer},
DeleteBucketRequest, DeleteBucketResponse, DeleteRequest, DeleteResponse, DeleteVersionsRequest, DeleteVersionsResponse, DeleteBucketRequest, DeleteBucketResponse, DeletePathsRequest, DeletePathsResponse, DeleteRequest, DeleteResponse,
DeleteVolumeRequest, DeleteVolumeResponse, GetBucketInfoRequest, GetBucketInfoResponse, ListBucketRequest, DeleteVersionRequest, DeleteVersionResponse, DeleteVersionsRequest, DeleteVersionsResponse, DeleteVolumeRequest,
ListBucketResponse, ListDirRequest, ListDirResponse, ListVolumesRequest, ListVolumesResponse, MakeBucketRequest, DeleteVolumeResponse, GetBucketInfoRequest, GetBucketInfoResponse, ListBucketRequest, ListBucketResponse, ListDirRequest,
MakeBucketResponse, MakeVolumeRequest, MakeVolumeResponse, MakeVolumesRequest, MakeVolumesResponse, PingRequest, ListDirResponse, ListVolumesRequest, ListVolumesResponse, MakeBucketRequest, MakeBucketResponse, MakeVolumeRequest,
PingResponse, ReadAllRequest, ReadAllResponse, ReadAtRequest, ReadAtResponse, ReadMultipleRequest, ReadMultipleResponse, MakeVolumeResponse, MakeVolumesRequest, MakeVolumesResponse, PingRequest, PingResponse, ReadAllRequest, ReadAllResponse,
ReadVersionRequest, ReadVersionResponse, ReadXlRequest, ReadXlResponse, RenameDataRequest, RenameDataResponse, ReadAtRequest, ReadAtResponse, ReadMultipleRequest, ReadMultipleResponse, ReadVersionRequest, ReadVersionResponse,
RenameFileRequst, RenameFileResponse, StatVolumeRequest, StatVolumeResponse, WalkDirRequest, WalkDirResponse, ReadXlRequest, ReadXlResponse, RenameDataRequest, RenameDataResponse, RenameFileRequst, RenameFileResponse,
WriteAllRequest, WriteAllResponse, WriteMetadataRequest, WriteMetadataResponse, WriteRequest, WriteResponse, RenamePartRequst, RenamePartResponse, StatVolumeRequest, StatVolumeResponse, UpdateMetadataRequest,
UpdateMetadataResponse, WalkDirRequest, WalkDirResponse, WriteAllRequest, WriteAllResponse, WriteMetadataRequest,
WriteMetadataResponse, WriteRequest, WriteResponse,
}, },
}; };
@@ -281,6 +283,36 @@ impl Node for NodeService {
} }
} }
async fn rename_part(&self, request: Request<RenamePartRequst>) -> Result<Response<RenamePartResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
match disk
.rename_part(
&request.src_volume,
&request.src_path,
&request.dst_volume,
&request.dst_path,
request.meta,
)
.await
{
Ok(_) => Ok(tonic::Response::new(RenamePartResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(RenamePartResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(RenamePartResponse {
success: false,
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn rename_file(&self, request: Request<RenameFileRequst>) -> Result<Response<RenameFileResponse>, Status> { async fn rename_file(&self, request: Request<RenameFileRequst>) -> Result<Response<RenameFileResponse>, Status> {
let request = request.into_inner(); let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await { if let Some(disk) = self.find_disk(&request.disk).await {
@@ -594,6 +626,68 @@ impl Node for NodeService {
} }
} }
async fn delete_paths(&self, request: Request<DeletePathsRequest>) -> Result<Response<DeletePathsResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let paths = request.paths.iter().map(|s| s.as_str()).collect::<Vec<&str>>();
match disk.delete_paths(&request.volume, &paths).await {
Ok(_) => Ok(tonic::Response::new(DeletePathsResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(DeletePathsResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(DeletePathsResponse {
success: false,
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn update_metadata(&self, request: Request<UpdateMetadataRequest>) -> Result<Response<UpdateMetadataResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let file_info = match serde_json::from_str::<FileInfo>(&request.file_info) {
Ok(file_info) => file_info,
Err(_) => {
return Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some("can not decode FileInfoVersions".to_string()),
}));
}
};
let opts = match serde_json::from_str::<UpdateMetadataOpts>(&request.opts) {
Ok(opts) => opts,
Err(_) => {
return Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some("can not decode UpdateMetadataOpts".to_string()),
}));
}
};
match disk.update_metadata(&request.volume, &request.path, file_info, opts).await {
Ok(_) => Ok(tonic::Response::new(UpdateMetadataResponse {
success: true,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(UpdateMetadataResponse {
success: false,
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn write_metadata(&self, request: Request<WriteMetadataRequest>) -> Result<Response<WriteMetadataResponse>, Status> { async fn write_metadata(&self, request: Request<WriteMetadataRequest>) -> Result<Response<WriteMetadataResponse>, Status> {
let request = request.into_inner(); let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await { if let Some(disk) = self.find_disk(&request.disk).await {
@@ -699,6 +793,60 @@ impl Node for NodeService {
} }
} }
async fn delete_version(&self, request: Request<DeleteVersionRequest>) -> Result<Response<DeleteVersionResponse>, Status> {
let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await {
let file_info = match serde_json::from_str::<FileInfo>(&request.file_info) {
Ok(file_info) => file_info,
Err(_) => {
return Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some("can not decode FileInfoVersions".to_string()),
}));
}
};
let opts = match serde_json::from_str::<DeleteOptions>(&request.opts) {
Ok(opts) => opts,
Err(_) => {
return Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some("can not decode DeleteOptions".to_string()),
}));
}
};
match disk
.delete_version(&request.volume, &request.path, file_info, request.force_del_marker, opts)
.await
{
Ok(raw_file_info) => match serde_json::to_string(&raw_file_info) {
Ok(raw_file_info) => Ok(tonic::Response::new(DeleteVersionResponse {
success: true,
raw_file_info,
error_info: None,
})),
Err(err) => Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some(err.to_string()),
})),
},
Err(err) => Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some(err.to_string()),
})),
}
} else {
Ok(tonic::Response::new(DeleteVersionResponse {
success: false,
raw_file_info: "".to_string(),
error_info: Some("can not find disk".to_string()),
}))
}
}
async fn delete_versions(&self, request: Request<DeleteVersionsRequest>) -> Result<Response<DeleteVersionsResponse>, Status> { async fn delete_versions(&self, request: Request<DeleteVersionsRequest>) -> Result<Response<DeleteVersionsResponse>, Status> {
let request = request.into_inner(); let request = request.into_inner();
if let Some(disk) = self.find_disk(&request.disk).await { if let Some(disk) = self.find_disk(&request.disk).await {
+2 -3
View File
@@ -19,7 +19,7 @@ use s3s::{auth::SimpleAuth, service::S3ServiceBuilder};
use service::hybrid; use service::hybrid;
use std::{io::IsTerminal, net::SocketAddr, str::FromStr}; use std::{io::IsTerminal, net::SocketAddr, str::FromStr};
use tokio::net::TcpListener; use tokio::net::TcpListener;
use tracing::{debug, info, warn}; use tracing::{debug, info};
use tracing_error::ErrorLayer; use tracing_error::ErrorLayer;
use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt}; use tracing_subscriber::{fmt, layer::SubscriberExt, util::SubscriberInitExt};
@@ -164,10 +164,9 @@ async fn run(opt: config::Opt) -> Result<()> {
} }
}); });
warn!(" init store");
// init store // init store
ECStore::new(opt.address.clone(), endpoint_pools.clone()).await?; ECStore::new(opt.address.clone(), endpoint_pools.clone()).await?;
warn!(" init store success!"); info!(" init store success!");
tokio::select! { tokio::select! {
_ = tokio::signal::ctrl_c() => { _ = tokio::signal::ctrl_c() => {
+1
View File
@@ -629,6 +629,7 @@ impl S3 for FS {
.. ..
} = req.input; } = req.input;
// error!("complete_multipart_upload {:?}", multipart_upload);
// mc cp step 5 // mc cp step 5
let Some(multipart_upload) = multipart_upload else { return Err(s3_error!(InvalidPart)) }; let Some(multipart_upload) = multipart_upload else { return Err(s3_error!(InvalidPart)) };
+12 -9
View File
@@ -1,24 +1,27 @@
#!/bin/bash #!/bin/bash
current_dir=$(pwd)
mkdir -p ./target/volume/test mkdir -p ./target/volume/test
mkdir -p ./target/volume/test{0..4} mkdir -p ./target/volume/test{0..4}
DATA_DIR="./target/volume/test"
# DATA_DIR="./target/volume/test{0...4}"
if [ -n "$1" ]; then
DATA_DIR="$1"
fi
if [ -z "$RUST_LOG" ]; then if [ -z "$RUST_LOG" ]; then
export RUST_LOG="rustfs=debug,ecstore=debug,s3s=debug" export RUST_LOG="rustfs=debug,ecstore=debug,s3s=debug"
fi fi
cargo run "$DATA_DIR" DATA_DIR_ARG="./target/volume/test{0...4}"
if [ -n "$1" ]; then
DATA_DIR_ARG="$1"
fi
# cargo run "$DATA_DIR_ARG"
# -- --access-key AKEXAMPLERUSTFS \ # -- --access-key AKEXAMPLERUSTFS \
# --secret-key SKEXAMPLERUSTFS \ # --secret-key SKEXAMPLERUSTFS \
# --address 0.0.0.0:9010 \ # --address 0.0.0.0:9010 \
# --domain-name 127.0.0.1:9010 \ # --domain-name 127.0.0.1:9010 \
"$DATA_DIR" # "$DATA_DIR_ARG"
# cargo run "$DATA_DIR" cargo run "$DATA_DIR_ARG"