From 33dc6c53a4d6103b6d65d2d145fd014c6b2fd723 Mon Sep 17 00:00:00 2001 From: weisd Date: Mon, 2 Dec 2024 00:02:31 +0800 Subject: [PATCH] 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)?;