mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 08:38:58 +00:00
fix: add set_disk
This commit is contained in:
@@ -10,6 +10,7 @@ pub mod error;
|
||||
mod file_meta;
|
||||
mod format;
|
||||
mod peer;
|
||||
pub mod set_disk;
|
||||
mod sets;
|
||||
pub mod store;
|
||||
pub mod store_api;
|
||||
|
||||
+22
-286
@@ -9,6 +9,7 @@ use crate::{
|
||||
endpoint::PoolEndpoints,
|
||||
erasure::Erasure,
|
||||
format::{DistributionAlgoVersion, FormatV3},
|
||||
set_disk::SetDisks,
|
||||
store_api::{
|
||||
BucketInfo, BucketOptions, FileInfo, MakeBucketOptions, MultipartUploadResult, ObjectOptions, PartInfo, PutObjReader,
|
||||
StorageAPI,
|
||||
@@ -23,7 +24,8 @@ use crate::{
|
||||
pub struct Sets {
|
||||
pub id: Uuid,
|
||||
// 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 pool_idx: usize,
|
||||
pub endpoints: PoolEndpoints,
|
||||
pub format: FormatV3,
|
||||
@@ -58,7 +60,15 @@ impl Sets {
|
||||
}
|
||||
}
|
||||
|
||||
disk_set.push(set_drive);
|
||||
let set_disks = SetDisks {
|
||||
disks: set_drive,
|
||||
set_drive_count,
|
||||
parity_count: partiy_count,
|
||||
set_index: i,
|
||||
pool_index: pool_idx,
|
||||
};
|
||||
|
||||
disk_set.push(set_disks);
|
||||
}
|
||||
|
||||
let sets = Self {
|
||||
@@ -76,11 +86,11 @@ impl Sets {
|
||||
|
||||
Ok(sets)
|
||||
}
|
||||
pub fn get_disks(&self, set_idx: usize) -> Vec<Option<DiskStore>> {
|
||||
pub fn get_disks(&self, set_idx: usize) -> SetDisks {
|
||||
self.disk_set[set_idx].clone()
|
||||
}
|
||||
|
||||
pub fn get_disks_by_key(&self, key: &str) -> Vec<Option<DiskStore>> {
|
||||
pub fn get_disks_by_key(&self, key: &str) -> SetDisks {
|
||||
self.get_disks(self.get_hashed_set_index(key))
|
||||
}
|
||||
|
||||
@@ -94,43 +104,6 @@ impl Sets {
|
||||
}
|
||||
}
|
||||
|
||||
async fn rename_data(
|
||||
&self,
|
||||
disks: &Vec<Option<DiskStore>>,
|
||||
src_bucket: &str,
|
||||
src_object: &str,
|
||||
file_infos: &Vec<FileInfo>,
|
||||
dst_bucket: &str,
|
||||
dst_object: &str,
|
||||
// write_quorum: usize,
|
||||
) -> Vec<Option<Error>> {
|
||||
let mut futures = Vec::with_capacity(disks.len());
|
||||
|
||||
for (i, disk) in disks.iter().enumerate() {
|
||||
let disk = disk.as_ref().unwrap();
|
||||
let file_info = &file_infos[i];
|
||||
futures.push(async move {
|
||||
disk.rename_data(src_bucket, src_object, file_info, dst_bucket, dst_object)
|
||||
.await
|
||||
})
|
||||
}
|
||||
|
||||
let mut errors = Vec::with_capacity(disks.len());
|
||||
|
||||
let results = join_all(futures).await;
|
||||
for result in results {
|
||||
match result {
|
||||
Ok(_) => {
|
||||
errors.push(None);
|
||||
}
|
||||
Err(e) => {
|
||||
errors.push(Some(e));
|
||||
}
|
||||
}
|
||||
}
|
||||
errors
|
||||
}
|
||||
|
||||
// async fn commit_rename_data_dir(
|
||||
// &self,
|
||||
// disks: &Vec<Option<DiskStore>>,
|
||||
@@ -143,61 +116,6 @@ impl Sets {
|
||||
// }
|
||||
}
|
||||
|
||||
async fn write_unique_file_info(
|
||||
disks: &Vec<Option<DiskStore>>,
|
||||
org_bucket: &str,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
files: &Vec<FileInfo>,
|
||||
// write_quorum: usize,
|
||||
) -> Vec<Option<Error>> {
|
||||
let mut futures = Vec::with_capacity(disks.len());
|
||||
|
||||
for (i, disk) in disks.iter().enumerate() {
|
||||
let disk = disk.as_ref().unwrap();
|
||||
let mut file_info = files[i].clone();
|
||||
file_info.erasure.index = i + 1;
|
||||
futures.push(async move { disk.write_metadata(org_bucket, bucket, prefix, file_info).await })
|
||||
}
|
||||
|
||||
let mut errors = Vec::with_capacity(disks.len());
|
||||
|
||||
let results = join_all(futures).await;
|
||||
for result in results {
|
||||
match result {
|
||||
Ok(_) => {
|
||||
errors.push(None);
|
||||
}
|
||||
Err(e) => {
|
||||
errors.push(Some(e));
|
||||
}
|
||||
}
|
||||
}
|
||||
errors
|
||||
}
|
||||
|
||||
fn get_upload_id_dir(bucket: &str, object: &str, upload_id: &str) -> String {
|
||||
let upload_uuid = match base64_decode(upload_id.as_bytes()) {
|
||||
Ok(res) => {
|
||||
let decoded_str = String::from_utf8(res).expect("Failed to convert decoded bytes to a UTF-8 string");
|
||||
let parts: Vec<&str> = decoded_str.splitn(2, '.').collect();
|
||||
if parts.len() == 2 {
|
||||
parts[1].to_string()
|
||||
} else {
|
||||
upload_id.to_string()
|
||||
}
|
||||
}
|
||||
Err(_) => upload_id.to_string(),
|
||||
};
|
||||
|
||||
format!("{}/{}", get_multipart_sha_dir(bucket, object), upload_uuid)
|
||||
}
|
||||
|
||||
fn get_multipart_sha_dir(bucket: &str, object: &str) -> String {
|
||||
let path = format!("{}/{}", bucket, object);
|
||||
hex(sha256(path.as_bytes()).as_ref())
|
||||
}
|
||||
|
||||
// #[derive(Debug)]
|
||||
// pub struct Objects {
|
||||
// pub endpoints: Vec<Endpoint>,
|
||||
@@ -219,107 +137,7 @@ impl StorageAPI for Sets {
|
||||
}
|
||||
|
||||
async fn put_object(&self, bucket: &str, object: &str, data: PutObjReader, opts: &ObjectOptions) -> Result<()> {
|
||||
let disks = self.get_disks_by_key(object);
|
||||
|
||||
let mut parity_drives = self.partiy_count;
|
||||
if opts.max_parity {
|
||||
parity_drives = disks.len() / 2;
|
||||
}
|
||||
|
||||
let data_drives = disks.len() - parity_drives;
|
||||
let mut write_quorum = data_drives;
|
||||
if data_drives == parity_drives {
|
||||
write_quorum += 1
|
||||
}
|
||||
|
||||
let mut fi = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives);
|
||||
|
||||
fi.data_dir = Uuid::new_v4();
|
||||
|
||||
let parts_metadata = vec![fi.clone(); disks.len()];
|
||||
|
||||
let (shuffle_disks, mut shuffle_parts_metadata) = shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi);
|
||||
|
||||
let mut writers = Vec::with_capacity(disks.len());
|
||||
|
||||
let mut futures = Vec::with_capacity(disks.len());
|
||||
|
||||
let tmp_dir = Uuid::new_v4().to_string();
|
||||
|
||||
let tmp_object = format!("{}/{}/part.1", tmp_dir, fi.data_dir);
|
||||
|
||||
for disk in shuffle_disks.iter() {
|
||||
let (reader, writer) = tokio::io::duplex(fi.erasure.block_size);
|
||||
|
||||
let disk = disk.as_ref().unwrap().clone();
|
||||
let tmp_object = tmp_object.clone();
|
||||
|
||||
// TODO: save small file in fileinfo.data instead of write file;
|
||||
|
||||
futures.push(async move {
|
||||
disk.create_file("", RUSTFS_META_TMP_BUCKET, tmp_object.as_str(), data.content_length, reader)
|
||||
.await
|
||||
});
|
||||
// futures.push(tokio::spawn(async move {
|
||||
// debug!("do createfile");
|
||||
// disk.CreateFile("", bucket.as_str(), object.as_str(), data.content_length, reader)
|
||||
// .await;
|
||||
// }));
|
||||
|
||||
writers.push(writer);
|
||||
}
|
||||
|
||||
let erasure = Erasure::new(fi.erasure.data_blocks, fi.erasure.parity_blocks);
|
||||
|
||||
let w_size = erasure
|
||||
.encode(data.stream, &mut writers, fi.erasure.block_size, data.content_length, write_quorum)
|
||||
.await?;
|
||||
|
||||
// close reader in create_file
|
||||
drop(writers);
|
||||
|
||||
let mut errors = Vec::with_capacity(disks.len());
|
||||
|
||||
let results = join_all(futures).await;
|
||||
for result in results {
|
||||
match result {
|
||||
Ok(_) => {
|
||||
errors.push(None);
|
||||
}
|
||||
Err(e) => {
|
||||
errors.push(Some(e));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
debug!("CreateFile errs:{:?}", errors);
|
||||
|
||||
// TODO: reduceWriteQuorumErrs
|
||||
// evalDisks
|
||||
|
||||
for fi in shuffle_parts_metadata.iter_mut() {
|
||||
fi.mod_time = OffsetDateTime::now_utc();
|
||||
fi.size = w_size;
|
||||
}
|
||||
|
||||
let rename_errs = self
|
||||
.rename_data(
|
||||
&shuffle_disks,
|
||||
RUSTFS_META_TMP_BUCKET,
|
||||
tmp_dir.as_str(),
|
||||
&shuffle_parts_metadata,
|
||||
&bucket,
|
||||
&object,
|
||||
)
|
||||
.await;
|
||||
|
||||
// TODO: reduceWriteQuorumErrs
|
||||
|
||||
debug!("put_object rename_errs:{:?}", rename_errs);
|
||||
|
||||
// self.commit_rename_data_dir(&shuffle_disks,&bucket,&object,)
|
||||
|
||||
Ok(())
|
||||
self.get_disks_by_key(object).put_object(bucket, object, data, opts).await
|
||||
}
|
||||
|
||||
async fn put_object_part(
|
||||
@@ -327,98 +145,16 @@ impl StorageAPI for Sets {
|
||||
bucket: &str,
|
||||
object: &str,
|
||||
upload_id: &str,
|
||||
_part_id: usize,
|
||||
_data: PutObjReader,
|
||||
_opts: &ObjectOptions,
|
||||
part_id: usize,
|
||||
data: PutObjReader,
|
||||
opts: &ObjectOptions,
|
||||
) -> Result<PartInfo> {
|
||||
let _upload_path = get_upload_id_dir(bucket, object, upload_id);
|
||||
|
||||
// TODO: checkUploadIDExists
|
||||
|
||||
unimplemented!()
|
||||
self.get_disks_by_key(object)
|
||||
.put_object_part(bucket, object, upload_id, part_id, data, opts)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn new_multipart_upload(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<MultipartUploadResult> {
|
||||
let disks = self.get_disks_by_key(object);
|
||||
|
||||
let mut parity_drives = self.partiy_count;
|
||||
if opts.max_parity {
|
||||
parity_drives = disks.len() / 2;
|
||||
}
|
||||
|
||||
let data_drives = disks.len() - parity_drives;
|
||||
let mut write_quorum = data_drives;
|
||||
if data_drives == parity_drives {
|
||||
write_quorum += 1
|
||||
}
|
||||
|
||||
let _ = write_quorum;
|
||||
|
||||
let mut fi = FileInfo::new([bucket, object].join("/").as_str(), data_drives, parity_drives);
|
||||
|
||||
fi.data_dir = Uuid::new_v4();
|
||||
fi.fresh = true;
|
||||
|
||||
let parts_metadata = vec![fi.clone(); disks.len()];
|
||||
|
||||
let (shuffle_disks, mut shuffle_parts_metadata) = shuffle_disks_and_parts_metadata(&disks, &parts_metadata, &fi);
|
||||
|
||||
for fi in shuffle_parts_metadata.iter_mut() {
|
||||
fi.mod_time = OffsetDateTime::now_utc();
|
||||
}
|
||||
|
||||
let upload_uuid = format!("{}x{}", Uuid::new_v4(), fi.mod_time);
|
||||
|
||||
let upload_id = base64_encode(format!("{}.{}", "globalDeploymentID", upload_uuid).as_bytes());
|
||||
|
||||
let upload_path = get_upload_id_dir(bucket, object, upload_uuid.as_str());
|
||||
|
||||
let errs = write_unique_file_info(
|
||||
&shuffle_disks,
|
||||
bucket,
|
||||
RUSTFS_META_MULTIPART_BUCKET,
|
||||
upload_path.as_str(),
|
||||
&shuffle_parts_metadata,
|
||||
)
|
||||
.await;
|
||||
|
||||
debug!("write_unique_file_info errs :{:?}", &errs);
|
||||
// TODO: reduceWriteQuorumErrs
|
||||
// evalDisks
|
||||
|
||||
Ok(MultipartUploadResult { upload_id })
|
||||
self.get_disks_by_key(object).new_multipart_upload(bucket, object, opts).await
|
||||
}
|
||||
}
|
||||
|
||||
// 打乱顺序
|
||||
fn shuffle_disks_and_parts_metadata(
|
||||
disks: &Vec<Option<DiskStore>>,
|
||||
parts_metadata: &Vec<FileInfo>,
|
||||
fi: &FileInfo,
|
||||
) -> (Vec<Option<DiskStore>>, Vec<FileInfo>) {
|
||||
let init = fi.mod_time == OffsetDateTime::UNIX_EPOCH;
|
||||
|
||||
let mut shuffled_disks = vec![None; disks.len()];
|
||||
let mut shuffled_parts_metadata = vec![FileInfo::default(); parts_metadata.len()];
|
||||
let distribution = &fi.erasure.distribution;
|
||||
|
||||
for (k, v) in disks.iter().enumerate() {
|
||||
if v.is_none() {
|
||||
continue;
|
||||
}
|
||||
|
||||
if !init && !parts_metadata[k].is_valid() {
|
||||
continue;
|
||||
}
|
||||
|
||||
// if !init && fi.xlv1 != parts_metadata[k].xlv1 {
|
||||
// continue;
|
||||
// }
|
||||
|
||||
let block_idx = distribution[k];
|
||||
shuffled_parts_metadata[block_idx - 1] = parts_metadata[k].clone();
|
||||
shuffled_disks[block_idx - 1] = disks[k].clone();
|
||||
}
|
||||
|
||||
(shuffled_disks, shuffled_parts_metadata)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user