diff --git a/.config/nextest.toml b/.config/nextest.toml index 06e69a29d..03529d03e 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -42,7 +42,7 @@ e2e-reliability = { max-threads = 1 } # --- default profile (local): serialize the flaky groups, never retry -------- [[profile.default.overrides]] -filter = 'package(rustfs-ecstore) & (test(concurrent_resend_same_part_commits_one_generation) | test(/^store::bucket::tests::bucket_delete_(mark_delete_marks|purge_removes|default_s3_delete)/))' +filter = 'package(rustfs-ecstore) & (test(concurrent_resend_same_part_commits_one_generation) | test(/^store::bucket::tests::bucket_delete_(mark_delete|purge_removes|default_s3_delete)/))' test-group = 'ecstore-serial-flaky' # Serialize the multipart crash-consistency scenarios (dist-2, backlog#1150): @@ -99,7 +99,7 @@ retries = 2 # QUARANTINE: OPEN backlog#937 — store::bucket::tests::bucket_delete_* race # make_bucket into InsufficientWriteQuorum via shared global state under load. [[profile.ci.overrides]] -filter = 'package(rustfs-ecstore) & test(/^store::bucket::tests::bucket_delete_(mark_delete_marks|purge_removes|default_s3_delete)/)' +filter = 'package(rustfs-ecstore) & test(/^store::bucket::tests::bucket_delete_(mark_delete|purge_removes|default_s3_delete)/)' test-group = 'ecstore-serial-flaky' retries = 2 diff --git a/.docker/test/site-replication/run-object-flow-check.sh b/.docker/test/site-replication/run-object-flow-check.sh index f351e2587..f96d69469 100755 --- a/.docker/test/site-replication/run-object-flow-check.sh +++ b/.docker/test/site-replication/run-object-flow-check.sh @@ -18,6 +18,7 @@ set -eu ACCESS_KEY="${RUSTFS_SITE_REPL_ACCESS_KEY:-rustfsadmin}" SECRET_KEY="${RUSTFS_SITE_REPL_SECRET_KEY:-rustfsadmin}" BUCKET="${RUSTFS_SITE_REPL_FLOW_BUCKET:-site-repl-flow-check}" +DELETE_BUCKET="${RUSTFS_SITE_REPL_DELETE_BUCKET:-site-repl-delete-$(date +%Y%m%d-%H%M%S)-$$}" PREFIX="${RUSTFS_SITE_REPL_FLOW_PREFIX:-flow-$(date +%Y%m%d-%H%M%S)}" WAIT_ATTEMPTS="${RUSTFS_SITE_REPL_WAIT_ATTEMPTS:-90}" WAIT_SLEEP_SECONDS="${RUSTFS_SITE_REPL_WAIT_SLEEP_SECONDS:-2}" @@ -85,17 +86,39 @@ wait_for_object() { wait_for_bucket() { site="$1" + bucket="${2:-$BUCKET}" attempt=1 while [ "$attempt" -le "$WAIT_ATTEMPTS" ]; do - if mc stat "$site/$BUCKET" >/dev/null 2>&1; then + if mc stat "$site/$bucket" >/dev/null 2>&1; then return 0 fi sleep "$WAIT_SLEEP_SECONDS" attempt=$((attempt + 1)) done - echo "bucket was not replicated in time: $site/$BUCKET" >&2 + echo "bucket was not replicated in time: $site/$bucket" >&2 + return 1 +} + +wait_for_bucket_delete() { + site="$1" + bucket="$2" + attempt=1 + + while [ "$attempt" -le "$WAIT_ATTEMPTS" ]; do + if result="$(mc stat --json "$site/$bucket" 2>&1)"; then + : + else + case "$result" in + *NoSuchBucket*) return 0 ;; + esac + fi + sleep "$WAIT_SLEEP_SECONDS" + attempt=$((attempt + 1)) + done + + echo "bucket deletion was not replicated in time: $site/$bucket" >&2 return 1 } @@ -186,6 +209,20 @@ EOF echo "verified replicated downloads for $object_name" done +echo "creating empty bucket for replicated delete check: $DELETE_BUCKET" +mc mb "site1/$DELETE_BUCKET" >/dev/null + +for site in site1 site2 site3; do + wait_for_bucket "$site" "$DELETE_BUCKET" +done + +echo "deleting empty bucket on site1: $DELETE_BUCKET" +mc rb "site1/$DELETE_BUCKET" >/dev/null + +for site in site1 site2 site3; do + wait_for_bucket_delete "$site" "$DELETE_BUCKET" +done + echo "site replication object flow check passed" echo "bucket: $BUCKET" echo "prefix: $PREFIX" diff --git a/.github/assets/site-replication-bucket-delete-epoch.png b/.github/assets/site-replication-bucket-delete-epoch.png new file mode 100644 index 000000000..9ed918779 Binary files /dev/null and b/.github/assets/site-replication-bucket-delete-epoch.png differ diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index c16588362..3f4b5602b 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -23,7 +23,7 @@ use crate::disk::{DiskAPI, DiskStore, disk_store::get_max_timeout_duration}; use crate::runtime::instance::{InstanceContext, bootstrap_ctx}; use crate::runtime::sources as runtime_sources; use crate::storage_api_contracts::bucket::{BucketInfo, BucketOptions, DeleteBucketOptions, MakeBucketOptions}; -use crate::store::utils::is_reserved_or_invalid_bucket; +use crate::store::{has_xlmeta_files, utils::is_reserved_or_invalid_bucket}; use crate::{ disk::{ self, VolumeInfo, @@ -614,13 +614,24 @@ impl PeerS3Client for LocalPeerS3Client { return Err(Error::ErasureWriteQuorum); } + let force = if opts.force_if_empty && !opts.force { + for disk in local_disks.iter() { + if has_xlmeta_files(&disk.path().join(bucket)).await.map_err(Error::Io)? { + return Err(Error::VolumeNotEmpty); + } + } + true + } else { + opts.force + }; + let mut futures = Vec::with_capacity(local_disks.len()); for disk in local_disks.iter() { // Non-force delete refuses a non-empty bucket (VolumeNotEmpty), which // the recreate loop below turns into BucketNotEmpty; only an explicit // force delete removes recursively (backlog#799 B1). - futures.push(disk.delete_volume(bucket, opts.force)); + futures.push(disk.delete_volume(bucket, force)); } let results = join_all(futures).await; @@ -988,13 +999,15 @@ impl PeerS3Client for RemotePeerS3Client { .await } - async fn delete_bucket(&self, bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { + async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { self.execute_with_timeout( || async { + let options = serde_json::to_string(opts)?; let mut client = self.get_client().await?; let request = Request::new(DeleteBucketRequest { bucket: bucket.to_string(), + options, }); let response = client.delete_bucket(request).await?.into_inner(); if !response.success { diff --git a/crates/ecstore/src/store/bucket.rs b/crates/ecstore/src/store/bucket.rs index a2968b3b5..0e25ea824 100644 --- a/crates/ecstore/src/store/bucket.rs +++ b/crates/ecstore/src/store/bucket.rs @@ -40,16 +40,16 @@ fn validate_table_bucket_delete_allowed( Ok(()) } -async fn table_catalog_metadata_exists(ctx: &crate::runtime::instance::InstanceContext, bucket: &str) -> bool { +async fn table_catalog_metadata_exists(ctx: &crate::runtime::instance::InstanceContext, bucket: &str) -> Result { let local_disks = runtime_sources::local_disks_in(ctx).await; for disk in local_disks.iter() { let catalog_path = disk.path().join(bucket).join(BUCKET_TABLE_RESERVED_PREFIX); - if has_xlmeta_files(&catalog_path).await { - return true; + if has_xlmeta_files(&catalog_path).await? { + return Ok(true); } } - false + Ok(false) } async fn validate_table_bucket_delete_guard(ctx: &crate::runtime::instance::InstanceContext, bucket: &str) -> Result<()> { @@ -57,7 +57,7 @@ async fn validate_table_bucket_delete_guard(ctx: &crate::runtime::instance::Inst .await .is_ok_and(|metadata| metadata.table_bucket_enabled()); if table_bucket_enabled { - validate_table_bucket_delete_allowed(bucket, true, table_catalog_metadata_exists(ctx, bucket).await)?; + validate_table_bucket_delete_allowed(bucket, true, table_catalog_metadata_exists(ctx, bucket).await?)?; } Ok(()) @@ -257,49 +257,57 @@ impl ECStore { None }; - // Check bucket exists before deletion (per S3 API spec) - // If bucket doesn't exist, return NoSuchBucket error - if let Err(err) = self.peer_sys.get_bucket_info(bucket, &BucketOptions::default()).await { - // Convert DiskError to StorageError for comparison - let storage_err: StorageError = err.into(); - if is_err_bucket_not_found(&storage_err) { - return Err(StorageError::BucketNotFound(bucket.to_string())); - } - return Err(to_object_err(storage_err, vec![bucket])); - } - - validate_table_bucket_delete_guard(&self.ctx, bucket).await?; - let sr_mark_delete = opts.srdelete_op == SRBucketDeleteOp::MarkDelete; let sr_purge = opts.srdelete_op == SRBucketDeleteOp::Purge; - - // Check bucket is empty before deletion (per S3 API spec) - // If bucket is not empty (contains actual objects with xl.meta files) and force - // is not set, return BucketNotEmpty error. - // Note: Empty directories (left after object deletion) should NOT count as objects. - if !opts.force && !sr_mark_delete { - let local_disks = runtime_sources::local_disks_in(&self.ctx).await; - for disk in local_disks.iter() { - // Check if bucket directory contains any xl.meta files (actual objects) - // We recursively scan for xl.meta files to determine if bucket has objects - // Use the disk's root path to construct bucket path - let bucket_path = disk.path().join(bucket); - if has_xlmeta_files(&bucket_path).await { - return Err(StorageError::BucketNotEmpty(bucket.to_string())); + let sr_delete = sr_mark_delete || sr_purge; + let mut delete_opts = opts.clone(); + let bucket_exists = match self.peer_sys.get_bucket_info(bucket, &BucketOptions::default()).await { + Ok(_) => true, + Err(err) => { + let storage_err: StorageError = err.into(); + if is_err_strict_volume_not_found(&storage_err) && sr_delete { + false + } else if is_err_strict_volume_not_found(&storage_err) { + return Err(StorageError::BucketNotFound(bucket.to_string())); + } else { + return Err(to_object_err(storage_err, vec![bucket])); } } + }; + + if bucket_exists { + validate_table_bucket_delete_guard(&self.ctx, bucket).await?; + + // Check bucket is empty before deletion (per S3 API spec) + // If bucket is not empty (contains actual objects with xl.meta files) and force + // is not set, return BucketNotEmpty error. + // Note: Empty directories (left after object deletion) should NOT count as objects. + if !opts.force { + let local_disks = runtime_sources::local_disks_in(&self.ctx).await; + for disk in local_disks.iter() { + let bucket_path = disk.path().join(bucket); + if has_xlmeta_files(&bucket_path).await? { + return Err(StorageError::BucketNotEmpty(bucket.to_string())); + } + } + delete_opts.force_if_empty = true; + } + } + + if sr_delete && !bucket_exists { + delete_opts.force_if_empty = true; } if sr_mark_delete { self.mark_bucket_deleted(bucket).await?; - self.cleanup_deleted_bucket_metadata(bucket, false).await?; - return Ok(()); } - self.peer_sys - .delete_bucket(bucket, opts) - .await - .map_err(|e| to_object_err(e.into(), vec![bucket]))?; + if let Err(err) = self.peer_sys.delete_bucket(bucket, &delete_opts).await { + let storage_err = to_object_err(err.into(), vec![bucket]); + if !sr_delete || !is_err_strict_volume_not_found(&storage_err) { + return Err(storage_err); + } + } self.cleanup_deleted_bucket_metadata(bucket, sr_purge).await?; Ok(()) @@ -319,7 +327,7 @@ mod tests { use crate::object_api::{ObjectOptions, PutObjReader}; use crate::runtime::instance::InstanceContext; use crate::storage_api_contracts::{ - bucket::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp}, + bucket::{BucketOperations as _, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp}, object::{ObjectIO as _, ObjectOperations as _}, }; use crate::store::{ECStore, init_local_disks_with_instance_ctx}; @@ -448,7 +456,10 @@ mod tests { async fn any_disk_has_object_metadata(disk_paths: &[PathBuf], bucket: &str) -> bool { for disk_path in disk_paths { - if super::has_xlmeta_files(&disk_path.join(bucket)).await { + if super::has_xlmeta_files(&disk_path.join(bucket)) + .await + .expect("object metadata scan should succeed") + { return true; } } @@ -508,12 +519,19 @@ mod tests { // serialize them so their assertions cannot observe each other's operations. #[tokio::test] #[serial] - async fn bucket_delete_mark_delete_marks_metadata_deleted_without_physical_object_delete() { + async fn bucket_delete_mark_delete_removes_empty_bucket_and_keeps_deleted_marker() { let (disk_paths, ecstore) = setup_bucket_delete_test_env().await; let bucket = format!("bucket-mark-delete-{}", Uuid::new_v4().simple()); - let object = "object.txt"; - create_bucket_with_object(&ecstore, &bucket, object).await; + ecstore + .make_bucket(&bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + for disk_path in &disk_paths { + tokio::fs::create_dir_all(disk_path.join(&bucket).join("empty-directory")) + .await + .expect("empty directory remnant should be created"); + } assert!(metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_ok()); let generation_before_delete = ecstore.scanner_namespace_mutation_generation(); @@ -526,7 +544,7 @@ mod tests { }, ) .await - .expect("MarkDelete should not reject non-empty bucket data"); + .expect("MarkDelete should remove an empty bucket"); assert_eq!( ecstore.scanner_namespace_mutation_generation(), generation_before_delete.saturating_add(1), @@ -534,8 +552,8 @@ mod tests { ); assert!( - any_disk_has_object_metadata(&disk_paths, &bucket).await, - "MarkDelete must not physically remove object xl.meta data" + !any_disk_path_exists(&disk_paths, &bucket).await, + "MarkDelete should remove the bucket volume" ); assert!( any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await, @@ -545,6 +563,79 @@ mod tests { metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_err(), "deleted bucket metadata must be removed from the local cache" ); + let buckets = ecstore + .list_bucket(&BucketOptions::default()) + .await + .expect("bucket listing should succeed after MarkDelete"); + assert!(!buckets.iter().any(|info| info.name == bucket)); + ecstore + .delete_all(RUSTFS_META_BUCKET, &bucket_deleted_marker_prefix(&bucket)) + .await + .expect("deleted-bucket marker should be removed to simulate a partial failure"); + assert!(!any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await); + ecstore + .delete_bucket( + &bucket, + &DeleteBucketOptions { + srdelete_op: SRBucketDeleteOp::MarkDelete, + ..Default::default() + }, + ) + .await + .expect("retried MarkDelete should recreate a missing tombstone"); + assert!(any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await); + } + + #[tokio::test] + #[serial] + async fn bucket_delete_mark_delete_rejects_non_empty_bucket_without_force() { + let (disk_paths, ecstore) = setup_bucket_delete_test_env().await; + let bucket = format!("bucket-mark-delete-non-empty-{}", Uuid::new_v4().simple()); + + create_bucket_with_object(&ecstore, &bucket, "object.txt").await; + + let err = ecstore + .delete_bucket( + &bucket, + &DeleteBucketOptions { + srdelete_op: SRBucketDeleteOp::MarkDelete, + ..Default::default() + }, + ) + .await + .expect_err("MarkDelete should reject a non-empty bucket without force"); + + assert!(matches!(err, StorageError::BucketNotEmpty(name) if name == bucket)); + assert!(any_disk_has_object_metadata(&disk_paths, &bucket).await); + assert!(!any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await); + } + + #[tokio::test] + #[serial] + async fn bucket_delete_mark_delete_rejects_hidden_object_paths_without_force() { + let (disk_paths, ecstore) = setup_bucket_delete_test_env().await; + let bucket = format!("bucket-mark-delete-hidden-{}", Uuid::new_v4().simple()); + + create_bucket_with_object(&ecstore, &bucket, ".well-known/acme-challenge").await; + let mut reader = PutObjReader::from_vec(b"delete bucket semantics".to_vec()); + ecstore + .put_object(&bucket, ".rustfs.sys/object", &mut reader, &ObjectOptions::default()) + .await + .expect("second hidden object should be written"); + + let err = ecstore + .delete_bucket( + &bucket, + &DeleteBucketOptions { + srdelete_op: SRBucketDeleteOp::MarkDelete, + ..Default::default() + }, + ) + .await + .expect_err("MarkDelete should reject a hidden object path without force"); + + assert!(matches!(err, StorageError::BucketNotEmpty(name) if name == bucket)); + assert!(any_disk_has_object_metadata(&disk_paths, &bucket).await); } #[tokio::test] @@ -557,7 +648,12 @@ mod tests { create_bucket_with_object(&ecstore, &bucket, object).await; write_bucket_metadata_marker(&disk_paths, &metadata_prefix).await; + ecstore + .mark_bucket_deleted(&bucket) + .await + .expect("deleted-bucket marker should be created"); assert!(any_disk_path_exists(&disk_paths, &metadata_prefix).await); + assert!(any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await); let generation_before_delete = ecstore.scanner_namespace_mutation_generation(); ecstore @@ -582,10 +678,32 @@ mod tests { !any_disk_path_exists(&disk_paths, &metadata_prefix).await, "Purge should remove bucket metadata prefix" ); + assert!( + !any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await, + "Purge should remove the deleted-bucket marker" + ); assert!( metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_err(), "purged bucket metadata must be removed from the local cache" ); + write_bucket_metadata_marker(&disk_paths, &metadata_prefix).await; + ecstore + .mark_bucket_deleted(&bucket) + .await + .expect("stale deleted-bucket marker should be recreated"); + ecstore + .delete_bucket( + &bucket, + &DeleteBucketOptions { + force: true, + srdelete_op: SRBucketDeleteOp::Purge, + ..Default::default() + }, + ) + .await + .expect("retried Purge should remove stale metadata without a bucket volume"); + assert!(!any_disk_path_exists(&disk_paths, &metadata_prefix).await); + assert!(!any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await); } #[tokio::test] diff --git a/crates/ecstore/src/store/mod.rs b/crates/ecstore/src/store/mod.rs index 90e851c17..4f8e8e7cb 100644 --- a/crates/ecstore/src/store/mod.rs +++ b/crates/ecstore/src/store/mod.rs @@ -40,8 +40,8 @@ use crate::disk::endpoint::{Endpoint, EndpointType}; use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions}; use crate::error::{Error, Result}; use crate::error::{ - StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_invalid_upload_id, is_err_object_not_found, - is_err_read_quorum, is_err_version_not_found, to_object_err, + StorageError, is_err_bucket_exists, is_err_invalid_upload_id, is_err_object_not_found, is_err_read_quorum, + is_err_strict_volume_not_found, is_err_version_not_found, to_object_err, }; use crate::runtime::global::DISK_RESERVE_FRACTION; use crate::runtime::instance::InstanceContext; @@ -91,7 +91,7 @@ type WalkOptions = StorageWalkOptions bool>; /// Check if a directory contains any xl.meta files (indicating actual S3 objects) /// This is used to determine if a bucket is empty for deletion purposes. -async fn has_xlmeta_files(path: &std::path::Path) -> bool { +pub(crate) async fn has_xlmeta_files(path: &std::path::Path) -> std::io::Result { use crate::disk::STORAGE_FORMAT_FILE; use tokio::fs; @@ -100,33 +100,27 @@ async fn has_xlmeta_files(path: &std::path::Path) -> bool { while let Some(current_path) = stack.pop() { let mut entries = match fs::read_dir(¤t_path).await { Ok(entries) => entries, - Err(_) => continue, + Err(err) if err.kind() == std::io::ErrorKind::NotFound => continue, + Err(err) => return Err(err), }; - while let Ok(Some(entry)) = entries.next_entry().await { + while let Some(entry) = entries.next_entry().await? { let file_name = entry.file_name(); let file_name_str = file_name.to_string_lossy(); - // Skip hidden files/directories (like .rustfs.sys) - if file_name_str.starts_with('.') { - continue; - } - // Check if this is an xl.meta file if file_name_str == STORAGE_FORMAT_FILE { - return true; + return Ok(true); } // If it's a directory, add to stack for further exploration - if let Ok(file_type) = entry.file_type().await - && file_type.is_dir() - { + if entry.file_type().await?.is_dir() { stack.push(entry.path()); } } } - false + Ok(false) } async fn enqueue_transition_after_write(result: Result, src: LcEventSrc) -> Result { diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index d22104b1d..0c3b8cb1d 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -83,6 +83,8 @@ pub struct GetBucketInfoResponse { pub struct DeleteBucketRequest { #[prost(string, tag = "1")] pub bucket: ::prost::alloc::string::String, + #[prost(string, tag = "2")] + pub options: ::prost::alloc::string::String, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct DeleteBucketResponse { diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 918f92c79..92b4e74a7 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -74,6 +74,7 @@ message GetBucketInfoResponse { message DeleteBucketRequest { string bucket = 1; + string options = 2; } message DeleteBucketResponse { diff --git a/crates/storage-api/src/bucket.rs b/crates/storage-api/src/bucket.rs index 1774fa8ad..d2e0411af 100644 --- a/crates/storage-api/src/bucket.rs +++ b/crates/storage-api/src/bucket.rs @@ -33,7 +33,7 @@ pub struct MakeBucketOptions { } /// Operation to perform on a bucket during site replication delete. -#[derive(Debug, Default, Clone, PartialEq)] +#[derive(Debug, Default, Clone, PartialEq, Serialize, Deserialize)] pub enum SRBucketDeleteOp { /// No operation. #[default] @@ -45,7 +45,7 @@ pub enum SRBucketDeleteOp { } /// Options for deleting a bucket. -#[derive(Debug, Default, Clone)] +#[derive(Debug, Default, Clone, Serialize, Deserialize)] pub struct DeleteBucketOptions { /// Skip acquiring namespace lock. pub no_lock: bool, @@ -53,6 +53,8 @@ pub struct DeleteBucketOptions { pub no_recreate: bool, /// Force deletion even if bucket is not empty. pub force: bool, + /// Force deletion only after the local peer verifies the bucket is empty. + pub force_if_empty: bool, /// Site replication delete operation. pub srdelete_op: SRBucketDeleteOp, } diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index 513745ef7..864fff455 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -2671,6 +2671,7 @@ mod tests { let request = Request::new(DeleteBucketRequest { bucket: "test-bucket".to_string(), + options: String::new(), }); let response = service.delete_bucket(request).await; @@ -2681,6 +2682,24 @@ mod tests { assert!(delete_response.success || delete_response.error.is_some()); } + #[tokio::test] + async fn test_delete_bucket_rejects_invalid_options() { + let service = create_test_node_service(); + + let request = Request::new(DeleteBucketRequest { + bucket: "test-bucket".to_string(), + options: "invalid json".to_string(), + }); + + let response = service + .delete_bucket(request) + .await + .expect("RPC response should be returned") + .into_inner(); + assert!(!response.success); + assert!(response.error.is_some()); + } + #[tokio::test] async fn test_read_all_invalid_disk() { let service = create_test_node_service(); diff --git a/rustfs/src/storage/rpc/node_service/bucket.rs b/rustfs/src/storage/rpc/node_service/bucket.rs index 3e7c76586..2d383f65b 100644 --- a/rustfs/src/storage/rpc/node_service/bucket.rs +++ b/rustfs/src/storage/rpc/node_service/bucket.rs @@ -103,17 +103,20 @@ impl NodeService { debug!("delete bucket"); let request = request.into_inner(); - match self - .local_peer - .delete_bucket( - &request.bucket, - &DeleteBucketOptions { - force: false, - ..Default::default() - }, - ) - .await - { + let options = if request.options.is_empty() { + DeleteBucketOptions::default() + } else { + match serde_json::from_str::(&request.options) { + Ok(options) => options, + Err(err) => { + return Ok(Response::new(DeleteBucketResponse { + success: false, + error: Some(DiskError::other(format!("decode DeleteBucketOptions failed: {err}")).into()), + })); + } + } + }; + match self.local_peer.delete_bucket(&request.bucket, &options).await { Ok(_) => Ok(Response::new(DeleteBucketResponse { success: true, error: None,