From bfd58e5749434ea738a6196cfa44b9e63644d963 Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 23 Dec 2024 23:24:11 +0800 Subject: [PATCH] fix:#175 add list_object_versions --- ecstore/src/disk/mod.rs | 118 ++++++++++++++++++++++++++++-- ecstore/src/set_disk.rs | 2 +- ecstore/src/sets.rs | 2 +- ecstore/src/store.rs | 36 ++------- ecstore/src/store_api.rs | 6 +- ecstore/src/store_list_objects.rs | 100 ++++++++++++++++++++++++- rustfs/src/storage/ecfs.rs | 73 +++++++++++++++++- 7 files changed, 292 insertions(+), 45 deletions(-) diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index c5f4fd9a9..254a2c9ff 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -558,6 +558,20 @@ pub struct FileInfoVersions { pub free_versions: Vec, } +impl FileInfoVersions { + pub fn find_version_index(&self, v: &str) -> Option { + 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> { + pub fn to_fileinfo(&self, bucket: &str) -> Result { 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 { @@ -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 { + 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)] diff --git a/ecstore/src/set_disk.rs b/ecstore/src/set_disk.rs index 962c9a7c4..59c77eaf8 100644 --- a/ecstore/src/set_disk.rs +++ b/ecstore/src/set_disk.rs @@ -3844,7 +3844,7 @@ impl StorageAPI for SetDisks { unimplemented!() } async fn list_object_versions( - &self, + self: Arc, _bucket: &str, _prefix: &str, _marker: &str, diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index ae095d2cc..8955ef4c2 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -463,7 +463,7 @@ impl StorageAPI for Sets { unimplemented!() } async fn list_object_versions( - &self, + self: Arc, _bucket: &str, _prefix: &str, _marker: &str, diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 1b41b247d..0aa548832 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -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, + bucket: &str, + prefix: &str, marker: &str, version_marker: &str, - _delimiter: &str, - _max_keys: i32, + delimiter: &str, + max_keys: i32, ) -> Result { - 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 { check_object_args(bucket, object)?; diff --git a/ecstore/src/store_api.rs b/ecstore/src/store_api.rs index 7dbd881ab..1cbcc487f 100644 --- a/ecstore/src/store_api.rs +++ b/ecstore/src/store_api.rs @@ -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, + pub next_version_idmarker: Option, pub objects: Vec, pub prefixes: Vec, } @@ -854,7 +854,7 @@ pub trait StorageAPI: ObjectIO { ) -> Result; // ListObjectVersions TODO: FIXME: async fn list_object_versions( - &self, + self: Arc, bucket: &str, prefix: &str, marker: &str, diff --git a/ecstore/src/store_list_objects.rs b/ecstore/src/store_list_objects.rs index f7b185d84..02c3e2ed8 100644 --- a/ecstore/src/store_list_objects.rs +++ b/ecstore/src/store_list_objects.rs @@ -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, + bucket: &str, + prefix: &str, + marker: &str, + version_marker: &str, + delimiter: &str, + max_keys: i32, + ) -> Result { + 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 = 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, o: &ListPathOptions) -> Result { check_list_objs_args(&o.bucket, &o.prefix, &o.marker)?; // if opts.prefix.ends_with(SLASH_SEPARATOR) { diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 2f808ef9a..da9eb186e 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -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, + req: S3Request, ) -> S3Result> { - 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 = 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))]