From 5594b189123bcc066951219b1a5ca5153c48150b Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Tue, 7 Jul 2026 08:25:17 +0800 Subject: [PATCH] fix(ecstore): make delete_volume non-recursive by default to prevent bucket-heal wipe (backlog#799 B1) (#4339) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(ecstore): make delete_volume non-recursive by default to prevent bucket-heal wipe (backlog#799 B1) `delete_volume` unconditionally `remove_dir_all`'d the whole bucket tree, and the bucket-heal "remove" branch called it fire-and-forget on every local disk. A mis-classified "dangling" bucket (or a non-force S3 DeleteBucket on a populated bucket) was therefore recursively wiped — a potential whole-bucket data loss. The `VolumeNotEmpty` -> recreate/`BucketNotEmpty` handling already present in both delete_bucket paths was dead code because the primitive never refused. Add an explicit `force_delete` flag to `DiskAPI::delete_volume` and default the non-force path to a non-recursive `remove_dir` (rmdir), which fails atomically with `VolumeNotEmpty` if the bucket still holds any object data. Only an explicit force delete (S3 force bucket delete) removes recursively. Mirrors MinIO's `xlStorage.DeleteVol` (`Remove` vs `RemoveAll`). - Trait + all impls (local behavior, dispatch, disk_store, remote RPC) take the flag; the gRPC `DeleteVolumeRequest` gains a `force` field (proto3 default false → old peers get the safe non-recursive behavior on rolling upgrade). - Heal remove branch passes `false` and no longer discards the result: a `VolumeNotEmpty` refusal is logged (the bucket is not dangling) instead of wiping data. - Both `delete_bucket` paths pass `opts.force`, activating the previously-dead `VolumeNotEmpty` -> `BucketNotEmpty`/recreate handling (correct S3 semantics). Adds a regression test: non-force delete of a non-empty bucket returns VolumeNotEmpty and preserves the data; force delete removes it. Design converged by two independent expert reviews (MinIO-fidelity + defense-in-depth) referencing MinIO xl-storage.go. Refs backlog#799 (B1), issue rustfs/backlog#850. The safety expert's deeper hardening (typed capability instead of a bool, trash-instead-of-in-place for force, quorum re-verification of dangling) is noted on #850 as follow-up. * fix(ecstore): reword 'mis-classified' -> 'misclassified' to satisfy typos (backlog#799 B1) * fix(rustfs): thread force_delete through StorageDiskRpcExt::delete_volume + test literal (backlog#799 B1) --- .../ecstore/src/cluster/rpc/peer_s3_client.rs | 23 +++++-- crates/ecstore/src/cluster/rpc/remote_disk.rs | 3 +- crates/ecstore/src/disk/disk_store.rs | 4 +- crates/ecstore/src/disk/local.rs | 64 +++++++++++++++++-- crates/ecstore/src/disk/mod.rs | 13 ++-- crates/ecstore/src/set_disk/ops/bucket.rs | 7 +- .../src/generated/proto_gen/node_service.rs | 2 + crates/protos/src/node.proto | 4 ++ rustfs/src/storage/rpc/node_service.rs | 1 + rustfs/src/storage/rpc/node_service/disk.rs | 2 +- rustfs/src/storage/storage_api.rs | 6 +- 11 files changed, 105 insertions(+), 24 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs index a6520b084..2640c60b3 100644 --- a/crates/ecstore/src/cluster/rpc/peer_s3_client.rs +++ b/crates/ecstore/src/cluster/rpc/peer_s3_client.rs @@ -553,7 +553,7 @@ impl PeerS3Client for LocalPeerS3Client { .ok_or(Error::VolumeNotFound) } - async fn delete_bucket(&self, bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { + async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { let local_disks = self.local_disks_for_pools().await; if local_disks.is_empty() { return Err(Error::ErasureWriteQuorum); @@ -562,7 +562,10 @@ impl PeerS3Client for LocalPeerS3Client { let mut futures = Vec::with_capacity(local_disks.len()); for disk in local_disks.iter() { - futures.push(disk.delete_volume(bucket)); + // 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)); } let results = join_all(futures).await; @@ -1036,9 +1039,19 @@ pub(crate) async fn heal_bucket_local_on_disks( futures.push(async move { match disk { Some(disk) => { - info!("will call delete_volume, volume: {}", bucket); - let _ = disk.delete_volume(&bucket).await; - None + // Non-force: a bucket that still holds object data refuses + // deletion (VolumeNotEmpty) instead of being recursively + // wiped, so a misclassified "dangling" bucket cannot lose + // data (backlog#799 B1). Surface that refusal instead of + // discarding it — it signals the bucket is not dangling. + match disk.delete_volume(&bucket, false).await { + Ok(()) => None, + Err(Error::VolumeNotEmpty) => { + warn!("heal declined to remove non-empty bucket {bucket} (not dangling)"); + None + } + Err(e) => Some(e), + } } None => Some(Error::DiskNotFound), } diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index ec86440d3..acd94726f 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -1339,7 +1339,7 @@ impl DiskAPI for RemoteDisk { } #[tracing::instrument(skip(self))] - async fn delete_volume(&self, volume: &str) -> Result<()> { + async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()> { debug!( event = EVENT_REMOTE_DISK_RPC, component = LOG_COMPONENT_ECSTORE, @@ -1360,6 +1360,7 @@ impl DiskAPI for RemoteDisk { let request = Request::new(DeleteVolumeRequest { disk: self.endpoint.to_string(), volume: volume.to_string(), + force: force_delete, }); let response = client.delete_volume(request).await?.into_inner(); diff --git a/crates/ecstore/src/disk/disk_store.rs b/crates/ecstore/src/disk/disk_store.rs index 8d179ffb9..c984718b8 100644 --- a/crates/ecstore/src/disk/disk_store.rs +++ b/crates/ecstore/src/disk/disk_store.rs @@ -1212,8 +1212,8 @@ impl DiskAPI for LocalDiskWrapper { .await } - async fn delete_volume(&self, volume: &str) -> Result<()> { - self.track_disk_health(|| async { self.disk.delete_volume(volume).await }, Duration::ZERO) + async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()> { + self.track_disk_health(|| async { self.disk.delete_volume(volume, force_delete).await }, Duration::ZERO) .await } diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 7bf9649ec..a0a8c50d3 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -4364,12 +4364,23 @@ impl DiskAPI for LocalDisk { } #[tracing::instrument(skip(self))] - async fn delete_volume(&self, volume: &str) -> Result<()> { + async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()> { let p = self.get_bucket_path(volume)?; - // TODO: avoid recursive deletion; return errVolumeNotEmpty when files remain + // Non-force is non-recursive: `remove_dir` (rmdir) fails atomically with + // `DirectoryNotEmpty` -> VolumeNotEmpty if the bucket still holds any + // object data, so a misclassified "dangling" bucket on the heal path + // (or a non-force S3 DeleteBucket on a populated bucket) can never be + // recursively wiped. Only an explicit `force_delete` (e.g. S3 force + // bucket delete) removes recursively. Mirrors MinIO's + // xlStorage.DeleteVol (Remove vs RemoveAll). (backlog#799 B1) + let res = if force_delete { + fs::remove_dir_all(&p).await + } else { + fs::remove_dir(&p).await + }; - if let Err(err) = fs::remove_dir_all(&p).await { + if let Err(err) = res { let e: DiskError = to_volume_error(err).into(); if e != DiskError::VolumeNotFound { return Err(e); @@ -5762,7 +5773,7 @@ mod test { disk.make_volumes(volumes.clone()).await.expect("operation should succeed"); - disk.delete_volume("a").await.expect("operation should succeed"); + disk.delete_volume("a", true).await.expect("operation should succeed"); let _ = fs::remove_dir_all(&p).await; } @@ -5871,7 +5882,48 @@ mod test { .expect("operation should succeed"); // Clean up - disk.delete_volume("test-volume").await.expect("operation should succeed"); + disk.delete_volume("test-volume", true) + .await + .expect("operation should succeed"); + let _ = fs::remove_dir_all(&test_dir).await; + } + + #[tokio::test] + async fn delete_volume_non_force_refuses_non_empty_bucket() { + // backlog#799 B1: a non-force delete_volume must refuse a bucket that + // still holds object data (VolumeNotEmpty) and leave it intact, so a + // misclassified "dangling" bucket cannot be recursively wiped. Only an + // explicit force delete removes it recursively. + let test_dir = "./test_b1_delete_volume_guard"; + let _ = fs::remove_dir_all(&test_dir).await; + fs::create_dir_all(&test_dir).await.expect("operation should succeed"); + let endpoint = Endpoint::try_from(test_dir).expect("operation should succeed"); + let disk = LocalDisk::new(&endpoint, false).await.expect("operation should succeed"); + + disk.make_volume("b1-bucket").await.expect("operation should succeed"); + let data: Vec = vec![1, 2, 3]; + disk.write_all("b1-bucket", "obj.dat", data.clone().into()) + .await + .expect("operation should succeed"); + + // Non-force must refuse and preserve the data. + let err = disk + .delete_volume("b1-bucket", false) + .await + .expect_err("non-empty bucket must be refused"); + assert!(matches!(err, DiskError::VolumeNotEmpty), "expected VolumeNotEmpty, got {err:?}"); + assert!( + disk.stat_volume("b1-bucket").await.is_ok(), + "bucket must still exist after a refused non-force delete" + ); + assert_eq!(disk.read_all("b1-bucket", "obj.dat").await.expect("data preserved"), data); + + // Force removes it recursively. + disk.delete_volume("b1-bucket", true) + .await + .expect("force delete removes non-empty"); + assert!(disk.stat_volume("b1-bucket").await.is_err(), "bucket must be gone after force delete"); + let _ = fs::remove_dir_all(&test_dir).await; } @@ -5899,7 +5951,7 @@ mod test { // Test deleting volumes for vol in &volumes { - disk.delete_volume(vol).await.expect("operation should succeed"); + disk.delete_volume(vol, true).await.expect("operation should succeed"); } // Clean up the test directory diff --git a/crates/ecstore/src/disk/mod.rs b/crates/ecstore/src/disk/mod.rs index abba9a357..af630dd16 100644 --- a/crates/ecstore/src/disk/mod.rs +++ b/crates/ecstore/src/disk/mod.rs @@ -187,10 +187,10 @@ impl DiskAPI for Disk { } #[tracing::instrument(skip(self))] - async fn delete_volume(&self, volume: &str) -> Result<()> { + async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()> { match self { - Disk::Local(local_disk) => local_disk.delete_volume(volume).await, - Disk::Remote(remote_disk) => remote_disk.delete_volume(volume).await, + Disk::Local(local_disk) => local_disk.delete_volume(volume, force_delete).await, + Disk::Remote(remote_disk) => remote_disk.delete_volume(volume, force_delete).await, } } @@ -554,7 +554,12 @@ pub trait DiskAPI: Debug + Send + Sync + 'static { async fn make_volumes(&self, volume: Vec<&str>) -> Result<()>; async fn list_volumes(&self) -> Result>; async fn stat_volume(&self, volume: &str) -> Result; - async fn delete_volume(&self, volume: &str) -> Result<()>; + /// Delete a volume (bucket directory). When `force_delete` is false a + /// non-empty volume is refused with `VolumeNotEmpty` (non-recursive); when + /// true it is removed recursively. Callers on the heal/dangling path must + /// pass false so a misclassified bucket that still holds data cannot be + /// wiped (backlog#799 B1). + async fn delete_volume(&self, volume: &str, force_delete: bool) -> Result<()>; // Concurrent read/write pipeline w <- MetaCacheEntry async fn walk_dir(&self, opts: WalkDirOptions, wr: &mut W) -> Result<()>; diff --git a/crates/ecstore/src/set_disk/ops/bucket.rs b/crates/ecstore/src/set_disk/ops/bucket.rs index 5b3f8393b..0d66aed89 100644 --- a/crates/ecstore/src/set_disk/ops/bucket.rs +++ b/crates/ecstore/src/set_disk/ops/bucket.rs @@ -179,16 +179,19 @@ impl BucketOperations for SetDisks { } #[tracing::instrument(skip(self))] - async fn delete_bucket(&self, bucket: &str, _opts: &DeleteBucketOptions) -> Result<()> { + async fn delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> { let disks = self.disk_inventory().await; let write_quorum = (disks.len() / 2) + 1; let mut futures = Vec::with_capacity(disks.len()); for disk in disks.iter().cloned() { let bucket = bucket.to_string(); + let force = opts.force; futures.push(async move { match disk { - Some(disk) => disk.delete_volume(&bucket).await, + // Non-force refuses a non-empty bucket (VolumeNotEmpty); only + // an explicit force delete removes recursively (backlog#799 B1). + Some(disk) => disk.delete_volume(&bucket, force).await, None => Err(DiskError::DiskNotFound), } }); diff --git a/crates/protos/src/generated/proto_gen/node_service.rs b/crates/protos/src/generated/proto_gen/node_service.rs index 68b095a5f..f381ffc80 100644 --- a/crates/protos/src/generated/proto_gen/node_service.rs +++ b/crates/protos/src/generated/proto_gen/node_service.rs @@ -650,6 +650,8 @@ pub struct DeleteVolumeRequest { pub disk: ::prost::alloc::string::String, #[prost(string, tag = "2")] pub volume: ::prost::alloc::string::String, + #[prost(bool, tag = "3")] + pub force: bool, } #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] pub struct DeleteVolumeResponse { diff --git a/crates/protos/src/node.proto b/crates/protos/src/node.proto index 5b667c8f6..a587b64ad 100644 --- a/crates/protos/src/node.proto +++ b/crates/protos/src/node.proto @@ -450,6 +450,10 @@ message ReadMultipleResponse { message DeleteVolumeRequest { string disk = 1; string volume = 2; + // When false (default), a non-empty volume is refused (VolumeNotEmpty); when + // true, the volume is deleted recursively. Old peers omit this field, which + // decodes to false — the safe, non-recursive behavior (backlog#799 B1). + bool force = 3; } message DeleteVolumeResponse { diff --git a/rustfs/src/storage/rpc/node_service.rs b/rustfs/src/storage/rpc/node_service.rs index aafdedd3d..102cc5ac5 100644 --- a/rustfs/src/storage/rpc/node_service.rs +++ b/rustfs/src/storage/rpc/node_service.rs @@ -2061,6 +2061,7 @@ mod tests { let request = Request::new(DeleteVolumeRequest { disk: "invalid-disk-path".to_string(), volume: "test-volume".to_string(), + force: false, }); let response = service.delete_volume(request).await; diff --git a/rustfs/src/storage/rpc/node_service/disk.rs b/rustfs/src/storage/rpc/node_service/disk.rs index de37815dc..9c640c481 100644 --- a/rustfs/src/storage/rpc/node_service/disk.rs +++ b/rustfs/src/storage/rpc/node_service/disk.rs @@ -159,7 +159,7 @@ impl NodeService { ) -> Result, Status> { let request = request.into_inner(); if let Some(disk) = self.find_disk(&request.disk).await { - match disk.delete_volume(&request.volume).await { + match disk.delete_volume(&request.volume, request.force).await { Ok(_) => Ok(Response::new(DeleteVolumeResponse { success: true, error: None, diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index a4f1ce103..a6304c66f 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -804,7 +804,7 @@ pub(crate) async fn update_erasure_type(setup_type: SetupType) { pub(crate) trait StorageDiskRpcExt { async fn disk_info(&self, opts: &DiskInfoOptions) -> DiskResult; - async fn delete_volume(&self, volume: &str) -> DiskResult<()>; + async fn delete_volume(&self, volume: &str, force_delete: bool) -> DiskResult<()>; async fn read_multiple(&self, req: ReadMultipleReq) -> DiskResult>; async fn batch_read_version(&self, req: BatchReadVersionReq) -> DiskResult>; async fn delete_versions(&self, volume: &str, versions: Vec, opts: DeleteOptions) @@ -884,8 +884,8 @@ where ecstore_disk::DiskAPI::disk_info(self, opts).await } - async fn delete_volume(&self, volume: &str) -> DiskResult<()> { - ecstore_disk::DiskAPI::delete_volume(self, volume).await + async fn delete_volume(&self, volume: &str, force_delete: bool) -> DiskResult<()> { + ecstore_disk::DiskAPI::delete_volume(self, volume, force_delete).await } async fn read_multiple(&self, req: ReadMultipleReq) -> DiskResult> {