mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-28 00:58:59 +00:00
fix:#175 add list_object_versions
This commit is contained in:
+110
-8
@@ -558,6 +558,20 @@ pub struct FileInfoVersions {
|
||||
pub free_versions: Vec<FileInfo>,
|
||||
}
|
||||
|
||||
impl FileInfoVersions {
|
||||
pub fn find_version_index(&self, v: &str) -> Option<usize> {
|
||||
if v.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let vid = Uuid::parse_str(v).unwrap_or(Uuid::nil());
|
||||
|
||||
for ver in self.versions.iter() {}
|
||||
|
||||
self.versions.iter().position(|v| v.version_id == Some(vid))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
||||
pub struct WalkDirOptions {
|
||||
// Bucket to scanner
|
||||
@@ -663,31 +677,31 @@ impl MetaCacheEntry {
|
||||
}
|
||||
|
||||
#[tracing::instrument(level = "debug", skip(self))]
|
||||
pub fn to_fileinfo(&self, bucket: &str) -> Result<Option<FileInfo>> {
|
||||
pub fn to_fileinfo(&self, bucket: &str) -> Result<FileInfo> {
|
||||
if self.is_dir() {
|
||||
return Ok(Some(FileInfo {
|
||||
return Ok(FileInfo {
|
||||
volume: bucket.to_owned(),
|
||||
name: self.name.clone(),
|
||||
..Default::default()
|
||||
}));
|
||||
});
|
||||
}
|
||||
|
||||
if self.cached.is_some() {
|
||||
let fm = self.cached.as_ref().unwrap();
|
||||
if fm.versions.is_empty() {
|
||||
return Ok(Some(FileInfo {
|
||||
return Ok(FileInfo {
|
||||
volume: bucket.to_owned(),
|
||||
name: self.name.clone(),
|
||||
deleted: true,
|
||||
is_latest: true,
|
||||
mod_time: Some(OffsetDateTime::UNIX_EPOCH),
|
||||
..Default::default()
|
||||
}));
|
||||
});
|
||||
}
|
||||
|
||||
let fi = fm.into_fileinfo(bucket, self.name.as_str(), "", false, false)?;
|
||||
|
||||
return Ok(Some(fi));
|
||||
return Ok(fi);
|
||||
}
|
||||
|
||||
let mut fm = FileMeta::new();
|
||||
@@ -695,7 +709,7 @@ impl MetaCacheEntry {
|
||||
|
||||
let fi = fm.into_fileinfo(bucket, self.name.as_str(), "", false, false)?;
|
||||
|
||||
return Ok(Some(fi));
|
||||
return Ok(fi);
|
||||
}
|
||||
|
||||
pub fn file_info_versions(&self, bucket: &str) -> Result<FileInfoVersions> {
|
||||
@@ -962,7 +976,7 @@ impl MetaCacheEntriesSorted {
|
||||
}
|
||||
}
|
||||
|
||||
if let Ok(Some(fi)) = entry.to_fileinfo(bucket) {
|
||||
if let Ok(fi) = entry.to_fileinfo(bucket) {
|
||||
// TODO:VersionPurgeStatus
|
||||
let versioned = vcfg.clone().map(|v| v.0.versioned(&entry.name)).unwrap_or_default();
|
||||
objects.push(fi.to_object_info(bucket, &entry.name, versioned));
|
||||
@@ -997,6 +1011,94 @@ impl MetaCacheEntriesSorted {
|
||||
|
||||
objects
|
||||
}
|
||||
|
||||
pub async fn file_info_versions(&self, bucket: &str, prefix: &str, delimiter: &str, after_v: &str) -> Vec<ObjectInfo> {
|
||||
let vcfg = get_versioning_config(bucket).await.ok();
|
||||
let mut objects = Vec::with_capacity(self.o.as_ref().len());
|
||||
let mut prev_prefix = "";
|
||||
let mut after_v = after_v;
|
||||
for entry in self.o.as_ref().iter().flatten() {
|
||||
if entry.is_object() {
|
||||
if !delimiter.is_empty() {
|
||||
if let Some(idx) = entry.name.trim_start_matches(prefix).find(delimiter) {
|
||||
let idx = prefix.len() + idx + delimiter.len();
|
||||
if let Some(curr_prefix) = entry.name.get(0..idx) {
|
||||
if curr_prefix == prev_prefix {
|
||||
continue;
|
||||
}
|
||||
|
||||
prev_prefix = curr_prefix;
|
||||
|
||||
objects.push(ObjectInfo {
|
||||
is_dir: true,
|
||||
bucket: bucket.to_owned(),
|
||||
name: curr_prefix.to_owned(),
|
||||
..Default::default()
|
||||
});
|
||||
}
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
let mut fiv = match entry.file_info_versions(bucket) {
|
||||
Ok(res) => res,
|
||||
Err(_err) => {
|
||||
//
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let fi_versions = 'c: {
|
||||
if !after_v.is_empty() {
|
||||
if let Some(idx) = fiv.find_version_index(after_v) {
|
||||
after_v = "";
|
||||
break 'c fiv.versions.split_off(idx + 1);
|
||||
}
|
||||
|
||||
after_v = "";
|
||||
break 'c fiv.versions;
|
||||
} else {
|
||||
break 'c fiv.versions;
|
||||
}
|
||||
};
|
||||
|
||||
for fi in fi_versions.into_iter() {
|
||||
// VersionPurgeStatus
|
||||
|
||||
let versioned = vcfg.clone().map(|v| v.0.versioned(&entry.name)).unwrap_or_default();
|
||||
objects.push(fi.to_object_info(bucket, &entry.name, versioned));
|
||||
}
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
if entry.is_dir() {
|
||||
if delimiter.is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Some(idx) = entry.name.trim_start_matches(prefix).find(delimiter) {
|
||||
let idx = prefix.len() + idx + delimiter.len();
|
||||
if let Some(curr_prefix) = entry.name.get(0..idx) {
|
||||
if curr_prefix == prev_prefix {
|
||||
continue;
|
||||
}
|
||||
|
||||
prev_prefix = curr_prefix;
|
||||
|
||||
objects.push(ObjectInfo {
|
||||
is_dir: true,
|
||||
bucket: bucket.to_owned(),
|
||||
name: curr_prefix.to_owned(),
|
||||
..Default::default()
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
objects
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Default)]
|
||||
|
||||
@@ -3844,7 +3844,7 @@ impl StorageAPI for SetDisks {
|
||||
unimplemented!()
|
||||
}
|
||||
async fn list_object_versions(
|
||||
&self,
|
||||
self: Arc<Self>,
|
||||
_bucket: &str,
|
||||
_prefix: &str,
|
||||
_marker: &str,
|
||||
|
||||
+1
-1
@@ -463,7 +463,7 @@ impl StorageAPI for Sets {
|
||||
unimplemented!()
|
||||
}
|
||||
async fn list_object_versions(
|
||||
&self,
|
||||
self: Arc<Self>,
|
||||
_bucket: &str,
|
||||
_prefix: &str,
|
||||
_marker: &str,
|
||||
|
||||
+8
-28
@@ -23,6 +23,7 @@ use crate::store_err::{
|
||||
to_object_err, StorageError,
|
||||
};
|
||||
use crate::store_init::ec_drives_no_config;
|
||||
use crate::store_list_objects::{max_keys_plus_one, ListPathOptions};
|
||||
use crate::utils::crypto::base64_decode;
|
||||
use crate::utils::path::{decode_dir_object, encode_dir_object, SLASH_SEPARATOR};
|
||||
use crate::utils::xml;
|
||||
@@ -1546,37 +1547,16 @@ impl StorageAPI for ECStore {
|
||||
// Ok(v2)
|
||||
}
|
||||
async fn list_object_versions(
|
||||
&self,
|
||||
_bucket: &str,
|
||||
_prefix: &str,
|
||||
self: Arc<Self>,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
marker: &str,
|
||||
version_marker: &str,
|
||||
_delimiter: &str,
|
||||
_max_keys: i32,
|
||||
delimiter: &str,
|
||||
max_keys: i32,
|
||||
) -> Result<ListObjectVersionsInfo> {
|
||||
if marker.is_empty() && !version_marker.is_empty() {
|
||||
return Err(Error::new(StorageError::NotImplemented));
|
||||
}
|
||||
|
||||
// let opts = ListPathOptions {
|
||||
// bucket: bucket.to_owned(),
|
||||
// marker: marker.to_owned(),
|
||||
// prefix: prefix.to_owned(),
|
||||
// limit: max_keys,
|
||||
// ..Default::default()
|
||||
// };
|
||||
|
||||
// let list = self
|
||||
// .list_path(&opts, delimiter)
|
||||
// .await
|
||||
// .map_err(|e| to_object_err(e, vec![bucket]))?;
|
||||
|
||||
// for info in list.objects.iter() {
|
||||
// //
|
||||
// }
|
||||
|
||||
// FIXME:
|
||||
unimplemented!()
|
||||
self.inner_list_object_versions(bucket, prefix, marker, version_marker, delimiter, max_keys)
|
||||
.await
|
||||
}
|
||||
async fn get_object_info(&self, bucket: &str, object: &str, opts: &ObjectOptions) -> Result<ObjectInfo> {
|
||||
check_object_args(bucket, object)?;
|
||||
|
||||
@@ -805,8 +805,8 @@ pub struct DeletedObject {
|
||||
|
||||
pub struct ListObjectVersionsInfo {
|
||||
pub is_truncated: bool,
|
||||
pub next_marker: String,
|
||||
pub next_version_idmarker: String,
|
||||
pub next_marker: Option<String>,
|
||||
pub next_version_idmarker: Option<String>,
|
||||
pub objects: Vec<ObjectInfo>,
|
||||
pub prefixes: Vec<String>,
|
||||
}
|
||||
@@ -854,7 +854,7 @@ pub trait StorageAPI: ObjectIO {
|
||||
) -> Result<ListObjectsV2Info>;
|
||||
// ListObjectVersions TODO: FIXME:
|
||||
async fn list_object_versions(
|
||||
&self,
|
||||
self: Arc<Self>,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
marker: &str,
|
||||
|
||||
@@ -6,7 +6,8 @@ use crate::file_meta::merge_file_meta_versions;
|
||||
use crate::peer::is_reserved_or_invalid_bucket;
|
||||
use crate::set_disk::SetDisks;
|
||||
use crate::store::check_list_objs_args;
|
||||
use crate::store_api::{ListObjectsInfo, ObjectInfo};
|
||||
use crate::store_api::{ListObjectVersionsInfo, ListObjectsInfo, ObjectInfo};
|
||||
use crate::store_err::StorageError;
|
||||
use crate::utils::path::{self, base_dir_from_prefix, SLASH_SEPARATOR};
|
||||
use crate::{store::ECStore, store_api::ListObjectsV2Info};
|
||||
use futures::future::join_all;
|
||||
@@ -26,7 +27,7 @@ const MAX_OBJECT_LIST: i32 = 1000;
|
||||
|
||||
const METACACHE_SHARE_PREFIX: bool = false;
|
||||
|
||||
fn max_keys_plus_one(max_keys: i32, add_one: bool) -> i32 {
|
||||
pub fn max_keys_plus_one(max_keys: i32, add_one: bool) -> i32 {
|
||||
let mut max_keys = max_keys;
|
||||
if max_keys > MAX_OBJECT_LIST {
|
||||
max_keys = MAX_OBJECT_LIST;
|
||||
@@ -227,6 +228,101 @@ impl ECStore {
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn inner_list_object_versions(
|
||||
self: Arc<Self>,
|
||||
bucket: &str,
|
||||
prefix: &str,
|
||||
marker: &str,
|
||||
version_marker: &str,
|
||||
delimiter: &str,
|
||||
max_keys: i32,
|
||||
) -> Result<ListObjectVersionsInfo> {
|
||||
if marker.is_empty() && !version_marker.is_empty() {
|
||||
return Err(Error::new(StorageError::NotImplemented));
|
||||
}
|
||||
|
||||
let opts = ListPathOptions {
|
||||
bucket: bucket.to_owned(),
|
||||
prefix: prefix.to_owned(),
|
||||
separator: delimiter.to_owned(),
|
||||
limit: max_keys_plus_one(max_keys, !marker.is_empty()),
|
||||
marker: marker.to_owned(),
|
||||
incl_deleted: true,
|
||||
ask_disks: "strict".to_owned(),
|
||||
versioned: true,
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let mut err_eof = false;
|
||||
let has_merged = match self.list_path(&opts).await {
|
||||
Ok(res) => Some(res),
|
||||
Err(err) => {
|
||||
if !is_err_eof(&err) {
|
||||
return Err(err);
|
||||
}
|
||||
|
||||
err_eof = true;
|
||||
None
|
||||
}
|
||||
};
|
||||
|
||||
let mut get_objects = if let Some(merged) = has_merged {
|
||||
merged.file_info_versions(bucket, prefix, delimiter, version_marker).await
|
||||
} else {
|
||||
Vec::new()
|
||||
};
|
||||
|
||||
let is_truncated = {
|
||||
if max_keys > 0 && get_objects.len() > max_keys as usize {
|
||||
get_objects.truncate(max_keys as usize);
|
||||
true
|
||||
} else {
|
||||
!err_eof && !get_objects.is_empty()
|
||||
}
|
||||
};
|
||||
|
||||
let mut prefixes: Vec<String> = Vec::new();
|
||||
|
||||
let mut objects = Vec::with_capacity(get_objects.len());
|
||||
for obj in get_objects.into_iter() {
|
||||
if obj.is_dir && obj.mod_time.is_none() && !delimiter.is_empty() {
|
||||
let mut found = false;
|
||||
if delimiter != SLASH_SEPARATOR {
|
||||
for p in prefixes.iter() {
|
||||
if found {
|
||||
break;
|
||||
}
|
||||
found = p == &obj.name;
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
prefixes.push(obj.name.clone());
|
||||
}
|
||||
} else {
|
||||
objects.push(obj);
|
||||
}
|
||||
}
|
||||
|
||||
let (next_marker, next_version_idmarker) = {
|
||||
if is_truncated {
|
||||
objects
|
||||
.last()
|
||||
.map(|last| (Some(last.name.clone()), last.version_id.map(|v| v.to_string())))
|
||||
.unwrap_or_default()
|
||||
} else {
|
||||
(None, None)
|
||||
}
|
||||
};
|
||||
|
||||
Ok(ListObjectVersionsInfo {
|
||||
is_truncated,
|
||||
next_marker,
|
||||
next_version_idmarker,
|
||||
objects,
|
||||
prefixes,
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn list_path(self: Arc<Self>, o: &ListPathOptions) -> Result<MetaCacheEntriesSorted> {
|
||||
check_list_objs_args(&o.bucket, &o.prefix, &o.marker)?;
|
||||
// if opts.prefix.ends_with(SLASH_SEPARATOR) {
|
||||
|
||||
@@ -503,6 +503,7 @@ impl S3 for FS {
|
||||
key: Some(v.name.to_owned()),
|
||||
last_modified: v.mod_time.map(Timestamp::from),
|
||||
size: Some(v.size as i64),
|
||||
e_tag: v.etag.clone(),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
@@ -541,9 +542,77 @@ impl S3 for FS {
|
||||
|
||||
async fn list_object_versions(
|
||||
&self,
|
||||
_req: S3Request<ListObjectVersionsInput>,
|
||||
req: S3Request<ListObjectVersionsInput>,
|
||||
) -> S3Result<S3Response<ListObjectVersionsOutput>> {
|
||||
Err(s3_error!(NotImplemented, "ListObjectVersions is not implemented yet"))
|
||||
let ListObjectVersionsInput {
|
||||
bucket,
|
||||
delimiter,
|
||||
key_marker,
|
||||
version_id_marker,
|
||||
max_keys,
|
||||
prefix,
|
||||
..
|
||||
} = req.input;
|
||||
|
||||
let prefix = prefix.unwrap_or_default();
|
||||
let delimiter = delimiter.unwrap_or_default();
|
||||
let max_keys = max_keys.unwrap_or(1000);
|
||||
|
||||
let Some(store) = new_object_layer_fn() else {
|
||||
return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string()));
|
||||
};
|
||||
|
||||
let object_infos = store
|
||||
.list_object_versions(
|
||||
&bucket,
|
||||
&prefix,
|
||||
&key_marker.unwrap_or_default(),
|
||||
&version_id_marker.unwrap_or_default(),
|
||||
&delimiter,
|
||||
max_keys,
|
||||
)
|
||||
.await
|
||||
.map_err(to_s3_error)?;
|
||||
|
||||
let objects: Vec<ObjectVersion> = object_infos
|
||||
.objects
|
||||
.iter()
|
||||
.filter(|v| !v.name.is_empty())
|
||||
.map(|v| {
|
||||
let obj = ObjectVersion {
|
||||
key: Some(v.name.to_owned()),
|
||||
last_modified: v.mod_time.map(Timestamp::from),
|
||||
size: Some(v.size as i64),
|
||||
version_id: v.version_id.map(|v| v.to_string()),
|
||||
is_latest: Some(v.is_latest),
|
||||
e_tag: v.etag.clone(),
|
||||
..Default::default() // TODO: another fields
|
||||
};
|
||||
|
||||
obj
|
||||
})
|
||||
.collect();
|
||||
|
||||
let key_count = objects.len() as i32;
|
||||
|
||||
let common_prefixes = object_infos
|
||||
.prefixes
|
||||
.into_iter()
|
||||
.map(|v| CommonPrefix { prefix: Some(v) })
|
||||
.collect();
|
||||
|
||||
let output = ListObjectVersionsOutput {
|
||||
// is_truncated: Some(object_infos.is_truncated),
|
||||
max_keys: Some(key_count),
|
||||
delimiter: Some(delimiter),
|
||||
name: Some(bucket),
|
||||
prefix: Some(prefix),
|
||||
common_prefixes: Some(common_prefixes),
|
||||
versions: Some(objects),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
Ok(S3Response::new(output))
|
||||
}
|
||||
|
||||
#[tracing::instrument(level = "debug", skip(self, req))]
|
||||
|
||||
Reference in New Issue
Block a user