fix(site-replication): delete replicated buckets

This commit is contained in:
Jason Kossis
2026-07-21 13:22:41 -04:00
committed by GitHub
parent f1d2af698c
commit 9469dfa5b8
11 changed files with 270 additions and 81 deletions
+2 -2
View File
@@ -42,7 +42,7 @@ e2e-reliability = { max-threads = 1 }
# --- default profile (local): serialize the flaky groups, never retry -------- # --- default profile (local): serialize the flaky groups, never retry --------
[[profile.default.overrides]] [[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' test-group = 'ecstore-serial-flaky'
# Serialize the multipart crash-consistency scenarios (dist-2, backlog#1150): # 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 # QUARANTINE: OPEN backlog#937 — store::bucket::tests::bucket_delete_* race
# make_bucket into InsufficientWriteQuorum via shared global state under load. # make_bucket into InsufficientWriteQuorum via shared global state under load.
[[profile.ci.overrides]] [[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' test-group = 'ecstore-serial-flaky'
retries = 2 retries = 2
@@ -18,6 +18,7 @@ set -eu
ACCESS_KEY="${RUSTFS_SITE_REPL_ACCESS_KEY:-rustfsadmin}" ACCESS_KEY="${RUSTFS_SITE_REPL_ACCESS_KEY:-rustfsadmin}"
SECRET_KEY="${RUSTFS_SITE_REPL_SECRET_KEY:-rustfsadmin}" SECRET_KEY="${RUSTFS_SITE_REPL_SECRET_KEY:-rustfsadmin}"
BUCKET="${RUSTFS_SITE_REPL_FLOW_BUCKET:-site-repl-flow-check}" 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)}" PREFIX="${RUSTFS_SITE_REPL_FLOW_PREFIX:-flow-$(date +%Y%m%d-%H%M%S)}"
WAIT_ATTEMPTS="${RUSTFS_SITE_REPL_WAIT_ATTEMPTS:-90}" WAIT_ATTEMPTS="${RUSTFS_SITE_REPL_WAIT_ATTEMPTS:-90}"
WAIT_SLEEP_SECONDS="${RUSTFS_SITE_REPL_WAIT_SLEEP_SECONDS:-2}" WAIT_SLEEP_SECONDS="${RUSTFS_SITE_REPL_WAIT_SLEEP_SECONDS:-2}"
@@ -85,17 +86,39 @@ wait_for_object() {
wait_for_bucket() { wait_for_bucket() {
site="$1" site="$1"
bucket="${2:-$BUCKET}"
attempt=1 attempt=1
while [ "$attempt" -le "$WAIT_ATTEMPTS" ]; do 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 return 0
fi fi
sleep "$WAIT_SLEEP_SECONDS" sleep "$WAIT_SLEEP_SECONDS"
attempt=$((attempt + 1)) attempt=$((attempt + 1))
done 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 return 1
} }
@@ -186,6 +209,20 @@ EOF
echo "verified replicated downloads for $object_name" echo "verified replicated downloads for $object_name"
done 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 "site replication object flow check passed"
echo "bucket: $BUCKET" echo "bucket: $BUCKET"
echo "prefix: $PREFIX" echo "prefix: $PREFIX"
Binary file not shown.

After

Width:  |  Height:  |  Size: 105 KiB

@@ -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::instance::{InstanceContext, bootstrap_ctx};
use crate::runtime::sources as runtime_sources; use crate::runtime::sources as runtime_sources;
use crate::storage_api_contracts::bucket::{BucketInfo, BucketOptions, DeleteBucketOptions, MakeBucketOptions}; 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::{ use crate::{
disk::{ disk::{
self, VolumeInfo, self, VolumeInfo,
@@ -614,13 +614,24 @@ impl PeerS3Client for LocalPeerS3Client {
return Err(Error::ErasureWriteQuorum); 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()); let mut futures = Vec::with_capacity(local_disks.len());
for disk in local_disks.iter() { for disk in local_disks.iter() {
// Non-force delete refuses a non-empty bucket (VolumeNotEmpty), which // Non-force delete refuses a non-empty bucket (VolumeNotEmpty), which
// the recreate loop below turns into BucketNotEmpty; only an explicit // the recreate loop below turns into BucketNotEmpty; only an explicit
// force delete removes recursively (backlog#799 B1). // 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; let results = join_all(futures).await;
@@ -988,13 +999,15 @@ impl PeerS3Client for RemotePeerS3Client {
.await .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( self.execute_with_timeout(
|| async { || async {
let options = serde_json::to_string(opts)?;
let mut client = self.get_client().await?; let mut client = self.get_client().await?;
let request = Request::new(DeleteBucketRequest { let request = Request::new(DeleteBucketRequest {
bucket: bucket.to_string(), bucket: bucket.to_string(),
options,
}); });
let response = client.delete_bucket(request).await?.into_inner(); let response = client.delete_bucket(request).await?.into_inner();
if !response.success { if !response.success {
+164 -46
View File
@@ -40,16 +40,16 @@ fn validate_table_bucket_delete_allowed(
Ok(()) 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<bool> {
let local_disks = runtime_sources::local_disks_in(ctx).await; let local_disks = runtime_sources::local_disks_in(ctx).await;
for disk in local_disks.iter() { for disk in local_disks.iter() {
let catalog_path = disk.path().join(bucket).join(BUCKET_TABLE_RESERVED_PREFIX); let catalog_path = disk.path().join(bucket).join(BUCKET_TABLE_RESERVED_PREFIX);
if has_xlmeta_files(&catalog_path).await { if has_xlmeta_files(&catalog_path).await? {
return true; return Ok(true);
} }
} }
false Ok(false)
} }
async fn validate_table_bucket_delete_guard(ctx: &crate::runtime::instance::InstanceContext, bucket: &str) -> Result<()> { 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 .await
.is_ok_and(|metadata| metadata.table_bucket_enabled()); .is_ok_and(|metadata| metadata.table_bucket_enabled());
if 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(()) Ok(())
@@ -257,49 +257,57 @@ impl ECStore {
None 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_mark_delete = opts.srdelete_op == SRBucketDeleteOp::MarkDelete;
let sr_purge = opts.srdelete_op == SRBucketDeleteOp::Purge; let sr_purge = opts.srdelete_op == SRBucketDeleteOp::Purge;
let sr_delete = sr_mark_delete || sr_purge;
// Check bucket is empty before deletion (per S3 API spec) let mut delete_opts = opts.clone();
// If bucket is not empty (contains actual objects with xl.meta files) and force let bucket_exists = match self.peer_sys.get_bucket_info(bucket, &BucketOptions::default()).await {
// is not set, return BucketNotEmpty error. Ok(_) => true,
// Note: Empty directories (left after object deletion) should NOT count as objects. Err(err) => {
if !opts.force && !sr_mark_delete { let storage_err: StorageError = err.into();
let local_disks = runtime_sources::local_disks_in(&self.ctx).await; if is_err_strict_volume_not_found(&storage_err) && sr_delete {
for disk in local_disks.iter() { false
// Check if bucket directory contains any xl.meta files (actual objects) } else if is_err_strict_volume_not_found(&storage_err) {
// We recursively scan for xl.meta files to determine if bucket has objects return Err(StorageError::BucketNotFound(bucket.to_string()));
// Use the disk's root path to construct bucket path } else {
let bucket_path = disk.path().join(bucket); return Err(to_object_err(storage_err, vec![bucket]));
if has_xlmeta_files(&bucket_path).await {
return Err(StorageError::BucketNotEmpty(bucket.to_string()));
} }
} }
};
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 { if sr_mark_delete {
self.mark_bucket_deleted(bucket).await?; self.mark_bucket_deleted(bucket).await?;
self.cleanup_deleted_bucket_metadata(bucket, false).await?;
return Ok(());
} }
self.peer_sys if let Err(err) = self.peer_sys.delete_bucket(bucket, &delete_opts).await {
.delete_bucket(bucket, opts) let storage_err = to_object_err(err.into(), vec![bucket]);
.await if !sr_delete || !is_err_strict_volume_not_found(&storage_err) {
.map_err(|e| to_object_err(e.into(), vec![bucket]))?; return Err(storage_err);
}
}
self.cleanup_deleted_bucket_metadata(bucket, sr_purge).await?; self.cleanup_deleted_bucket_metadata(bucket, sr_purge).await?;
Ok(()) Ok(())
@@ -319,7 +327,7 @@ mod tests {
use crate::object_api::{ObjectOptions, PutObjReader}; use crate::object_api::{ObjectOptions, PutObjReader};
use crate::runtime::instance::InstanceContext; use crate::runtime::instance::InstanceContext;
use crate::storage_api_contracts::{ use crate::storage_api_contracts::{
bucket::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp}, bucket::{BucketOperations as _, BucketOptions, DeleteBucketOptions, MakeBucketOptions, SRBucketDeleteOp},
object::{ObjectIO as _, ObjectOperations as _}, object::{ObjectIO as _, ObjectOperations as _},
}; };
use crate::store::{ECStore, init_local_disks_with_instance_ctx}; 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 { async fn any_disk_has_object_metadata(disk_paths: &[PathBuf], bucket: &str) -> bool {
for disk_path in disk_paths { 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; return true;
} }
} }
@@ -508,12 +519,19 @@ mod tests {
// serialize them so their assertions cannot observe each other's operations. // serialize them so their assertions cannot observe each other's operations.
#[tokio::test] #[tokio::test]
#[serial] #[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 (disk_paths, ecstore) = setup_bucket_delete_test_env().await;
let bucket = format!("bucket-mark-delete-{}", Uuid::new_v4().simple()); 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()); assert!(metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_ok());
let generation_before_delete = ecstore.scanner_namespace_mutation_generation(); let generation_before_delete = ecstore.scanner_namespace_mutation_generation();
@@ -526,7 +544,7 @@ mod tests {
}, },
) )
.await .await
.expect("MarkDelete should not reject non-empty bucket data"); .expect("MarkDelete should remove an empty bucket");
assert_eq!( assert_eq!(
ecstore.scanner_namespace_mutation_generation(), ecstore.scanner_namespace_mutation_generation(),
generation_before_delete.saturating_add(1), generation_before_delete.saturating_add(1),
@@ -534,8 +552,8 @@ mod tests {
); );
assert!( assert!(
any_disk_has_object_metadata(&disk_paths, &bucket).await, !any_disk_path_exists(&disk_paths, &bucket).await,
"MarkDelete must not physically remove object xl.meta data" "MarkDelete should remove the bucket volume"
); );
assert!( assert!(
any_disk_path_exists(&disk_paths, bucket_deleted_marker_volume(&bucket)).await, 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(), metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_err(),
"deleted bucket metadata must be removed from the local cache" "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] #[tokio::test]
@@ -557,7 +648,12 @@ mod tests {
create_bucket_with_object(&ecstore, &bucket, object).await; create_bucket_with_object(&ecstore, &bucket, object).await;
write_bucket_metadata_marker(&disk_paths, &metadata_prefix).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, &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(); let generation_before_delete = ecstore.scanner_namespace_mutation_generation();
ecstore ecstore
@@ -582,10 +678,32 @@ mod tests {
!any_disk_path_exists(&disk_paths, &metadata_prefix).await, !any_disk_path_exists(&disk_paths, &metadata_prefix).await,
"Purge should remove bucket metadata prefix" "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!( assert!(
metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_err(), metadata_sys::get_in(&ecstore.ctx, &bucket).await.is_err(),
"purged bucket metadata must be removed from the local cache" "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] #[tokio::test]
+9 -15
View File
@@ -40,8 +40,8 @@ use crate::disk::endpoint::{Endpoint, EndpointType};
use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions}; use crate::disk::{DiskAPI, DiskInfo, DiskInfoOptions};
use crate::error::{Error, Result}; use crate::error::{Error, Result};
use crate::error::{ use crate::error::{
StorageError, is_err_bucket_exists, is_err_bucket_not_found, is_err_invalid_upload_id, is_err_object_not_found, StorageError, is_err_bucket_exists, is_err_invalid_upload_id, is_err_object_not_found, is_err_read_quorum,
is_err_read_quorum, is_err_version_not_found, to_object_err, is_err_strict_volume_not_found, is_err_version_not_found, to_object_err,
}; };
use crate::runtime::global::DISK_RESERVE_FRACTION; use crate::runtime::global::DISK_RESERVE_FRACTION;
use crate::runtime::instance::InstanceContext; use crate::runtime::instance::InstanceContext;
@@ -91,7 +91,7 @@ type WalkOptions = StorageWalkOptions<fn(&FileInfo) -> bool>;
/// Check if a directory contains any xl.meta files (indicating actual S3 objects) /// 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. /// 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<bool> {
use crate::disk::STORAGE_FORMAT_FILE; use crate::disk::STORAGE_FORMAT_FILE;
use tokio::fs; use tokio::fs;
@@ -100,33 +100,27 @@ async fn has_xlmeta_files(path: &std::path::Path) -> bool {
while let Some(current_path) = stack.pop() { while let Some(current_path) = stack.pop() {
let mut entries = match fs::read_dir(&current_path).await { let mut entries = match fs::read_dir(&current_path).await {
Ok(entries) => entries, 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 = entry.file_name();
let file_name_str = file_name.to_string_lossy(); 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 // Check if this is an xl.meta file
if file_name_str == STORAGE_FORMAT_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 it's a directory, add to stack for further exploration
if let Ok(file_type) = entry.file_type().await if entry.file_type().await?.is_dir() {
&& file_type.is_dir()
{
stack.push(entry.path()); stack.push(entry.path());
} }
} }
} }
false Ok(false)
} }
async fn enqueue_transition_after_write(result: Result<ObjectInfo>, src: LcEventSrc) -> Result<ObjectInfo> { async fn enqueue_transition_after_write(result: Result<ObjectInfo>, src: LcEventSrc) -> Result<ObjectInfo> {
@@ -83,6 +83,8 @@ pub struct GetBucketInfoResponse {
pub struct DeleteBucketRequest { pub struct DeleteBucketRequest {
#[prost(string, tag = "1")] #[prost(string, tag = "1")]
pub bucket: ::prost::alloc::string::String, pub bucket: ::prost::alloc::string::String,
#[prost(string, tag = "2")]
pub options: ::prost::alloc::string::String,
} }
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)] #[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct DeleteBucketResponse { pub struct DeleteBucketResponse {
+1
View File
@@ -74,6 +74,7 @@ message GetBucketInfoResponse {
message DeleteBucketRequest { message DeleteBucketRequest {
string bucket = 1; string bucket = 1;
string options = 2;
} }
message DeleteBucketResponse { message DeleteBucketResponse {
+4 -2
View File
@@ -33,7 +33,7 @@ pub struct MakeBucketOptions {
} }
/// Operation to perform on a bucket during site replication delete. /// 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 { pub enum SRBucketDeleteOp {
/// No operation. /// No operation.
#[default] #[default]
@@ -45,7 +45,7 @@ pub enum SRBucketDeleteOp {
} }
/// Options for deleting a bucket. /// Options for deleting a bucket.
#[derive(Debug, Default, Clone)] #[derive(Debug, Default, Clone, Serialize, Deserialize)]
pub struct DeleteBucketOptions { pub struct DeleteBucketOptions {
/// Skip acquiring namespace lock. /// Skip acquiring namespace lock.
pub no_lock: bool, pub no_lock: bool,
@@ -53,6 +53,8 @@ pub struct DeleteBucketOptions {
pub no_recreate: bool, pub no_recreate: bool,
/// Force deletion even if bucket is not empty. /// Force deletion even if bucket is not empty.
pub force: bool, pub force: bool,
/// Force deletion only after the local peer verifies the bucket is empty.
pub force_if_empty: bool,
/// Site replication delete operation. /// Site replication delete operation.
pub srdelete_op: SRBucketDeleteOp, pub srdelete_op: SRBucketDeleteOp,
} }
+19
View File
@@ -2671,6 +2671,7 @@ mod tests {
let request = Request::new(DeleteBucketRequest { let request = Request::new(DeleteBucketRequest {
bucket: "test-bucket".to_string(), bucket: "test-bucket".to_string(),
options: String::new(),
}); });
let response = service.delete_bucket(request).await; let response = service.delete_bucket(request).await;
@@ -2681,6 +2682,24 @@ mod tests {
assert!(delete_response.success || delete_response.error.is_some()); 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] #[tokio::test]
async fn test_read_all_invalid_disk() { async fn test_read_all_invalid_disk() {
let service = create_test_node_service(); let service = create_test_node_service();
+14 -11
View File
@@ -103,17 +103,20 @@ impl NodeService {
debug!("delete bucket"); debug!("delete bucket");
let request = request.into_inner(); let request = request.into_inner();
match self let options = if request.options.is_empty() {
.local_peer DeleteBucketOptions::default()
.delete_bucket( } else {
&request.bucket, match serde_json::from_str::<DeleteBucketOptions>(&request.options) {
&DeleteBucketOptions { Ok(options) => options,
force: false, Err(err) => {
..Default::default() return Ok(Response::new(DeleteBucketResponse {
}, success: false,
) error: Some(DiskError::other(format!("decode DeleteBucketOptions failed: {err}")).into()),
.await }));
{ }
}
};
match self.local_peer.delete_bucket(&request.bucket, &options).await {
Ok(_) => Ok(Response::new(DeleteBucketResponse { Ok(_) => Ok(Response::new(DeleteBucketResponse {
success: true, success: true,
error: None, error: None,