From df60aa50fedfa0c81d3711a675adff5ba9d4eb0f Mon Sep 17 00:00:00 2001 From: weisd Date: Sat, 30 Nov 2024 20:54:51 +0800 Subject: [PATCH 1/2] todo file_info_from_raw --- ecstore/src/bucket/versioning/mod.rs | 2 +- ecstore/src/disk/local.rs | 2 ++ ecstore/src/disk/mod.rs | 1 + ecstore/src/sets.rs | 2 ++ ecstore/src/store.rs | 1 + rustfs/src/storage/ecfs.rs | 13 +++++++++++-- 6 files changed, 18 insertions(+), 3 deletions(-) diff --git a/ecstore/src/bucket/versioning/mod.rs b/ecstore/src/bucket/versioning/mod.rs index 3444034e4..fb3387cc9 100644 --- a/ecstore/src/bucket/versioning/mod.rs +++ b/ecstore/src/bucket/versioning/mod.rs @@ -19,7 +19,7 @@ impl VersioningApi for VersioningConfiguration { } fn prefix_enabled(&self, prefix: &str) -> bool { - if self.status == Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)) { + if self.status != Some(BucketVersioningStatus::from_static(BucketVersioningStatus::ENABLED)) { return false; } diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 46dee48a6..3951d7899 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -399,6 +399,7 @@ impl LocalDisk { } /// read xl.meta raw data + #[tracing::instrument(level = "debug", skip(self, volume_dir, path))] async fn read_raw( &self, bucket: &str, @@ -460,6 +461,7 @@ impl LocalDisk { Ok(data) } + #[tracing::instrument(level = "debug", skip(self, volume_dir, file_path))] async fn read_all_data_with_dmtime( &self, volume: &str, diff --git a/ecstore/src/disk/mod.rs b/ecstore/src/disk/mod.rs index 10ceeafb7..6896b751e 100644 --- a/ecstore/src/disk/mod.rs +++ b/ecstore/src/disk/mod.rs @@ -277,6 +277,7 @@ impl DiskAPI for Disk { } } + #[tracing::instrument(level = "debug", skip(self))] async fn read_version( &self, _org_volume: &str, diff --git a/ecstore/src/sets.rs b/ecstore/src/sets.rs index 8dd190839..cc177b771 100644 --- a/ecstore/src/sets.rs +++ b/ecstore/src/sets.rs @@ -273,6 +273,7 @@ struct DelObj { #[async_trait::async_trait] impl ObjectIO for Sets { + #[tracing::instrument(level = "debug", skip(self))] async fn get_object_reader( &self, bucket: &str, @@ -285,6 +286,7 @@ impl ObjectIO for Sets { .get_object_reader(bucket, object, range, h, opts) .await } + #[tracing::instrument(level = "debug", skip(self, data))] async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result { self.get_disks_by_key(object).put_object(bucket, object, data, opts).await } diff --git a/ecstore/src/store.rs b/ecstore/src/store.rs index 9787a4ad9..cf6e9f5b5 100644 --- a/ecstore/src/store.rs +++ b/ecstore/src/store.rs @@ -1083,6 +1083,7 @@ impl ObjectIO for ECStore { .get_object_reader(bucket, object.as_str(), range, h, &opts) .await } + #[tracing::instrument(level = "debug", skip(self, data))] async fn put_object(&self, bucket: &str, object: &str, data: &mut PutObjReader, opts: &ObjectOptions) -> Result { check_put_object_args(bucket, object)?; diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 8241b2353..bc69569a1 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -303,14 +303,21 @@ impl S3 for FS { let range = HTTPRangeSpec::nil(); let h = HeaderMap::new(); - let opts = &ObjectOptions::default(); + + let metadata = extract_metadata(&req.headers); + + let opts: ObjectOptions = put_opts(&bucket, &key, None, &req.headers, Some(metadata)) + .await + .map_err(to_s3_error)?; + + error!("get_object ObjectOptions {:?}", &opts); let Some(store) = new_object_layer_fn() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; let reader = store - .get_object_reader(bucket.as_str(), key.as_str(), range, h, opts) + .get_object_reader(bucket.as_str(), key.as_str(), range, h, &opts) .await .map_err(to_s3_error)?; @@ -574,6 +581,8 @@ impl S3 for FS { .await .map_err(to_s3_error)?; + error!("ObjectOptions {:?}", opts); + let obj_info = store .put_object(&bucket, &key, &mut reader, &opts) .await From 33dc6c53a4d6103b6d65d2d145fd014c6b2fd723 Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 2 Dec 2024 00:02:31 +0800 Subject: [PATCH 2/2] fix download err where versioning enable --- ecstore/src/disk/local.rs | 1 + ecstore/src/options.rs | 46 ++++++++++++++++++++++++++++++++++++++ rustfs/src/storage/ecfs.rs | 32 +++++++++++++++++--------- 3 files changed, 69 insertions(+), 10 deletions(-) diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 3951d7899..cb05d22e3 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -1845,6 +1845,7 @@ impl DiskAPI for LocalDisk { self.delete_file(&volume_dir, &xl_path, true, false).await } + #[tracing::instrument(level = "debug", skip(self))] async fn delete_versions( &self, volume: &str, diff --git a/ecstore/src/options.rs b/ecstore/src/options.rs index 2da19e2c6..e2920e19d 100644 --- a/ecstore/src/options.rs +++ b/ecstore/src/options.rs @@ -8,6 +8,52 @@ use lazy_static::lazy_static; use std::collections::HashMap; use uuid::Uuid; +pub async fn del_opts( + bucket: &str, + object: &str, + vid: Option, + headers: &HeaderMap, + metadata: Option>, +) -> Result { + let versioned = BucketVersioningSys::prefix_enabled(bucket, object).await; + let version_suspended = BucketVersioningSys::prefix_suspended(bucket, object).await; + + let vid = vid.map(|v| v.as_str().trim().to_owned()); + + if let Some(ref id) = vid { + if let Err(_err) = Uuid::parse_str(id.as_str()) { + return Err(Error::new(StorageError::InvalidVersionID( + bucket.to_owned(), + object.to_owned(), + id.clone(), + ))); + } + + if !versioned { + return Err(Error::new(StorageError::InvalidArgument( + bucket.to_owned(), + object.to_owned(), + id.clone(), + ))); + } + } + + let mut opts = put_opts_from_headers(headers, metadata) + .map_err(|err| Error::new(StorageError::InvalidArgument(bucket.to_owned(), object.to_owned(), err.to_string())))?; + + opts.version_id = { + if is_dir_object(object) && vid.is_none() { + Some(Uuid::nil().to_string()) + } else { + vid + } + }; + opts.version_suspended = version_suspended; + opts.versioned = versioned; + + Ok(opts) +} + pub async fn put_opts( bucket: &str, object: &str, diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index bc69569a1..6000090b6 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -16,6 +16,7 @@ use ecstore::bucket::tagging::decode_tags; use ecstore::bucket::tagging::encode_tags; use ecstore::bucket::versioning_sys::BucketVersioningSys; use ecstore::new_object_layer_fn; +use ecstore::options::del_opts; use ecstore::options::extract_metadata; use ecstore::options::put_opts; use ecstore::store_api::BucketOptions; @@ -155,7 +156,14 @@ impl S3 for FS { bucket, key, version_id, .. } = req.input; - let version_id = version_id + let metadata = extract_metadata(&req.headers); + + let opts: ObjectOptions = del_opts(&bucket, &key, version_id, &req.headers, Some(metadata)) + .await + .map_err(to_s3_error)?; + + let version_id = opts + .version_id .as_ref() .map(|v| match Uuid::parse_str(v) { Ok(id) => Some(id), @@ -172,10 +180,7 @@ impl S3 for FS { let Some(store) = new_object_layer_fn() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - let (dobjs, _errs) = store - .delete_objects(&bucket, objects, ObjectOptions::default()) - .await - .map_err(to_s3_error)?; + let (dobjs, _errs) = store.delete_objects(&bucket, objects, opts).await.map_err(to_s3_error)?; // TODO: let errors; @@ -240,11 +245,13 @@ impl S3 for FS { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); }; - let (dobjs, _errs) = store - .delete_objects(&bucket, objects, ObjectOptions::default()) + let metadata = extract_metadata(&req.headers); + + let opts: ObjectOptions = del_opts(&bucket, "", None, &req.headers, Some(metadata)) .await .map_err(to_s3_error)?; - // info!("delete_objects res {:?} {:?}", &dobjs, errs); + + let (dobjs, errs) = store.delete_objects(&bucket, objects, opts).await.map_err(to_s3_error)?; let deleted = dobjs .iter() @@ -263,6 +270,9 @@ impl S3 for FS { .collect(); // TODO: let errors; + for err in errs.iter().flatten() { + warn!("delete_objects err {:?}", err); + } let output = DeleteObjectsOutput { deleted: Some(deleted), @@ -298,7 +308,9 @@ impl S3 for FS { async fn get_object(&self, req: S3Request) -> S3Result> { // mc get 3 - let GetObjectInput { bucket, key, .. } = req.input; + let GetObjectInput { + bucket, key, version_id, .. + } = req.input; let range = HTTPRangeSpec::nil(); @@ -306,7 +318,7 @@ impl S3 for FS { let metadata = extract_metadata(&req.headers); - let opts: ObjectOptions = put_opts(&bucket, &key, None, &req.headers, Some(metadata)) + let opts: ObjectOptions = put_opts(&bucket, &key, version_id, &req.headers, Some(metadata)) .await .map_err(to_s3_error)?;