rewrite real_meta

This commit is contained in:
weisd
2024-12-04 17:36:16 +08:00
committed by weisd
parent 4ccd69c782
commit 2786174ffb
6 changed files with 278 additions and 134 deletions
+104 -103
View File
@@ -24,19 +24,19 @@ use crate::store_err::{
};
use crate::store_init::ec_drives_no_config;
use crate::utils::crypto::base64_decode;
use crate::utils::path::{base_dir_from_prefix, decode_dir_object, encode_dir_object, SLASH_SEPARATOR};
use crate::utils::path::{decode_dir_object, encode_dir_object, SLASH_SEPARATOR};
use crate::utils::xml;
use crate::{
bucket::metadata::BucketMetadata,
disk::{error::DiskError, new_disk, DiskOption, DiskStore, WalkDirOptions, BUCKET_META_PREFIX, RUSTFS_META_BUCKET},
disk::{error::DiskError, new_disk, DiskOption, DiskStore, BUCKET_META_PREFIX, RUSTFS_META_BUCKET},
endpoints::EndpointServerPools,
error::{Error, Result},
peer::S3PeerSys,
sets::Sets,
store_api::{
BucketInfo, BucketOptions, CompletePart, DeleteBucketOptions, DeletedObject, GetObjectReader, HTTPRangeSpec,
ListObjectsInfo, ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, ObjectToDelete,
PartInfo, PutObjReader, StorageAPI,
ListObjectsV2Info, MakeBucketOptions, MultipartUploadResult, ObjectInfo, ObjectOptions, ObjectToDelete, PartInfo,
PutObjReader, StorageAPI,
},
store_init, utils,
};
@@ -52,11 +52,7 @@ use std::cmp::Ordering;
use std::process::exit;
use std::slice::Iter;
use std::time::SystemTime;
use std::{
collections::{HashMap, HashSet},
sync::Arc,
time::Duration,
};
use std::{collections::HashMap, sync::Arc, time::Duration};
use time::OffsetDateTime;
use tokio::select;
use tokio::sync::mpsc::Sender;
@@ -265,102 +261,104 @@ impl ECStore {
self.pools.len() == 1
}
pub async fn list_path(&self, opts: &ListPathOptions, delimiter: &str) -> Result<ListObjectsInfo> {
// if opts.prefix.ends_with(SLASH_SEPARATOR) {
// return Err(Error::msg("eof"));
// }
// define in store_list_objects.rs
// pub async fn list_path(&self, opts: &ListPathOptions, delimiter: &str) -> Result<ListObjectsInfo> {
// // if opts.prefix.ends_with(SLASH_SEPARATOR) {
// // return Err(Error::msg("eof"));
// // }
let mut opts = opts.clone();
// let mut opts = opts.clone();
if opts.base_dir.is_empty() {
opts.base_dir = base_dir_from_prefix(&opts.prefix);
}
// if opts.base_dir.is_empty() {
// opts.base_dir = base_dir_from_prefix(&opts.prefix);
// }
let objects = self.list_merged(&opts, delimiter).await?;
// let objects = self.list_merged(&opts, delimiter).await?;
let info = ListObjectsInfo {
objects,
..Default::default()
};
Ok(info)
}
// let info = ListObjectsInfo {
// objects,
// ..Default::default()
// };
// Ok(info)
// }
// 读所有
async fn list_merged(&self, opts: &ListPathOptions, delimiter: &str) -> Result<Vec<ObjectInfo>> {
let walk_opts = WalkDirOptions {
bucket: opts.bucket.clone(),
base_dir: opts.base_dir.clone(),
..Default::default()
};
// define in store_list_objects.rs
// async fn list_merged(&self, opts: &ListPathOptions, delimiter: &str) -> Result<Vec<ObjectInfo>> {
// let walk_opts = WalkDirOptions {
// bucket: opts.bucket.clone(),
// base_dir: opts.base_dir.clone(),
// ..Default::default()
// };
// let (mut wr, mut rd) = tokio::io::duplex(1024);
// // let (mut wr, mut rd) = tokio::io::duplex(1024);
let mut futures = Vec::new();
// let mut futures = Vec::new();
for sets in self.pools.iter() {
for set in sets.disk_set.iter() {
futures.push(set.walk_dir(&walk_opts));
}
}
// for sets in self.pools.iter() {
// for set in sets.disk_set.iter() {
// futures.push(set.walk_dir(&walk_opts));
// }
// }
let results = join_all(futures).await;
// let results = join_all(futures).await;
// let mut errs = Vec::new();
let mut ress = Vec::new();
let mut uniq = HashSet::new();
// // let mut errs = Vec::new();
// let mut ress = Vec::new();
// let mut uniq = HashSet::new();
for (disks_ress, _disks_errs) in results {
for disks_res in disks_ress.iter() {
if disks_res.is_none() {
// TODO handle errs
continue;
}
let entrys = disks_res.as_ref().unwrap();
// for (disks_ress, _disks_errs) in results {
// for disks_res in disks_ress.iter() {
// if disks_res.is_none() {
// // TODO handle errs
// continue;
// }
// let entrys = disks_res.as_ref().unwrap();
for entry in entrys {
// warn!("lst_merged entry---- {}", &entry.name);
// for entry in entrys {
// // warn!("lst_merged entry---- {}", &entry.name);
if !opts.prefix.is_empty() && !entry.name.starts_with(&opts.prefix) {
continue;
}
// if !opts.prefix.is_empty() && !entry.name.starts_with(&opts.prefix) {
// continue;
// }
if !uniq.contains(&entry.name) {
uniq.insert(entry.name.clone());
// TODO: 过滤
// if !uniq.contains(&entry.name) {
// uniq.insert(entry.name.clone());
// // TODO: 过滤
if opts.limit > 0 && ress.len() as i32 >= opts.limit {
return Ok(ress);
}
// if opts.limit > 0 && ress.len() as i32 >= opts.limit {
// return Ok(ress);
// }
if entry.is_object() {
if !delimiter.is_empty() {
// entry.name.trim_start_matches(pat)
}
// if entry.is_object() {
// if !delimiter.is_empty() {
// // entry.name.trim_start_matches(pat)
// }
let fi = entry.to_fileinfo(&opts.bucket)?;
if let Some(f) = fi {
ress.push(f.to_object_info(&opts.bucket, &entry.name, false));
}
continue;
}
// let fi = entry.to_fileinfo(&opts.bucket)?;
// if let Some(f) = fi {
// ress.push(f.to_object_info(&opts.bucket, &entry.name, false));
// }
// continue;
// }
if entry.is_dir() {
ress.push(ObjectInfo {
is_dir: true,
bucket: opts.bucket.clone(),
name: entry.name.clone(),
..Default::default()
});
}
}
}
}
}
// if entry.is_dir() {
// ress.push(ObjectInfo {
// is_dir: true,
// bucket: opts.bucket.clone(),
// name: entry.name.clone(),
// ..Default::default()
// });
// }
// }
// }
// }
// }
// warn!("list_merged errs {:?}", errs);
// // warn!("list_merged errs {:?}", errs);
Ok(ress)
}
// Ok(ress)
// }
async fn delete_all(&self, bucket: &str, prefix: &str) -> Result<()> {
let mut futures = Vec::new();
@@ -1520,29 +1518,32 @@ impl StorageAPI for ECStore {
continuation_token: &str,
delimiter: &str,
max_keys: i32,
_fetch_owner: bool,
_start_after: &str,
fetch_owner: bool,
start_after: &str,
) -> Result<ListObjectsV2Info> {
let opts = ListPathOptions {
bucket: bucket.to_string(),
limit: max_keys,
prefix: prefix.to_owned(),
..Default::default()
};
self.inner_list_objects_v2(&bucket, &prefix, &continuation_token, &delimiter, max_keys, fetch_owner, start_after)
.await
let info = self.list_path(&opts, delimiter).await?;
// let opts = ListPathOptions {
// bucket: bucket.to_string(),
// limit: max_keys,
// prefix: prefix.to_owned(),
// ..Default::default()
// };
// warn!("list_objects_v2 info {:?}", info);
// let info = self.list_path(&opts, delimiter).await?;
let v2 = ListObjectsV2Info {
is_truncated: info.is_truncated,
continuation_token: continuation_token.to_owned(),
next_continuation_token: info.next_marker,
objects: info.objects,
prefixes: info.prefixes,
};
// // warn!("list_objects_v2 info {:?}", info);
Ok(v2)
// let v2 = ListObjectsV2Info {
// is_truncated: info.is_truncated,
// continuation_token: continuation_token.to_owned(),
// next_continuation_token: info.next_marker,
// objects: info.objects,
// prefixes: info.prefixes,
// };
// Ok(v2)
}
async fn list_object_versions(
&self,
@@ -2238,7 +2239,7 @@ fn check_bucket_and_object_names(bucket: &str, object: &str) -> Result<()> {
Ok(())
}
fn check_list_objs_args(bucket: &str, prefix: &str, _marker: &str) -> Result<()> {
pub fn check_list_objs_args(bucket: &str, prefix: &str, _marker: &str) -> Result<()> {
if !is_meta_bucketname(bucket) && check_valid_bucket_name_strict(bucket).is_err() {
return Err(Error::new(StorageError::BucketNameInvalid(bucket.to_string())));
}