From cfd6e2b32fbdefe473852a60911b6ab854f9e547 Mon Sep 17 00:00:00 2001 From: weisd Date: Tue, 27 Aug 2024 14:12:15 +0800 Subject: [PATCH] add: delete_versions --- Cargo.lock | 1 + ecstore/src/disk/local.rs | 4 ++++ rustfs/Cargo.toml | 2 +- rustfs/src/storage/ecfs.rs | 29 ++++++++++++++++++++++++++++- 4 files changed, 34 insertions(+), 2 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index dd42ff16a..8773cca8c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1168,6 +1168,7 @@ dependencies = [ "tracing-error", "tracing-subscriber", "transform-stream", + "uuid", ] [[package]] diff --git a/ecstore/src/disk/local.rs b/ecstore/src/disk/local.rs index 407c779e4..3daf2637c 100644 --- a/ecstore/src/disk/local.rs +++ b/ecstore/src/disk/local.rs @@ -292,6 +292,7 @@ impl LocalDisk { // 没有版本了,删除xl.meta if fm.versions.is_empty() { self.delete_file(&volume_dir, &xlpath, true, false).await?; + return Ok(()); } // 更新xl.meta @@ -871,6 +872,9 @@ impl DiskAPI for LocalDisk { _opts: DeleteOptions, ) -> Result>> { let mut errs = Vec::with_capacity(versions.len()); + for _ in 0..versions.len() { + errs.push(None); + } for (i, ver) in versions.iter().enumerate() { if let Err(e) = self.delete_versions_internal(volume, ver.name.as_str(), &ver.versions).await { diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 5e5414f0f..f08b977e0 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -23,7 +23,7 @@ http.workspace = true bytes.workspace = true futures.workspace = true futures-util.workspace = true - +uuid = { version = "1.8.0", features = ["v4", "fast-rng", "serde"] } ecstore = { path = "../ecstore" } s3s = "0.10.0" clap = { version = "4.5.7", features = ["derive"] } diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index 8e8f53db5..822488aa0 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -6,6 +6,7 @@ use ecstore::store_api::HTTPRangeSpec; use ecstore::store_api::MakeBucketOptions; use ecstore::store_api::MultipartUploadResult; use ecstore::store_api::ObjectOptions; +use ecstore::store_api::ObjectToDelete; use ecstore::store_api::PutObjReader; use ecstore::store_api::StorageAPI; use futures::pin_mut; @@ -20,7 +21,9 @@ use s3s::S3; use s3s::{S3Request, S3Response}; use std::fmt::Debug; use std::str::FromStr; +use tracing::info; use transform_stream::AsyncTryStream; +use uuid::Uuid; use ecstore::error::Result; use ecstore::store::ECStore; @@ -99,7 +102,31 @@ impl S3 for FS { #[tracing::instrument(level = "debug", skip(self, req))] async fn delete_objects(&self, req: S3Request) -> S3Result> { - let _input = req.input; + info!("delete_objects args {:?}", req.input); + + let DeleteObjectsInput { bucket, delete, .. } = req.input; + + let objects: Vec = delete + .objects + .iter() + .map(|v| { + let version_id = v + .version_id + .as_ref() + .map(|v| match Uuid::parse_str(v) { + Ok(id) => Some(id), + Err(_) => None, + }) + .unwrap_or_default(); + ObjectToDelete { + object_name: v.key.clone(), + version_id: version_id, + } + }) + .collect(); + + let (dobjs, errs) = try_!(self.store.delete_objects(&bucket, objects, ObjectOptions::default()).await); + info!("delete_objects res {:?} {:?}", &dobjs, errs); let output = DeleteObjectsOutput { ..Default::default() }; Ok(S3Response::new(output))