mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-07 20:46:11 +00:00
fix(ecstore): defer zero-evidence delete diagnostics (#7036)
This commit is contained in:
@@ -198,6 +198,11 @@ async fn bucket_delete_local_blocker(
|
|||||||
}
|
}
|
||||||
if scan.diagnostic_truncated {
|
if scan.diagnostic_truncated {
|
||||||
residue.diagnostic_truncated = true;
|
residue.diagnostic_truncated = true;
|
||||||
|
if residue.files == 0 {
|
||||||
|
// A zero-evidence timeout is inconclusive. Preflight still uses
|
||||||
|
// non-recursive `force_if_empty`; post-failure keeps the physical error.
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
record_bucket_delete_blocker(bucket, BucketDeleteBlockerKind::DiagnosticBudgetExceeded, &residue);
|
record_bucket_delete_blocker(bucket, BucketDeleteBlockerKind::DiagnosticBudgetExceeded, &residue);
|
||||||
return Ok(Some(StorageError::BucketNotEmptyWithDetails {
|
return Ok(Some(StorageError::BucketNotEmptyWithDetails {
|
||||||
bucket: bucket.to_string(),
|
bucket: bucket.to_string(),
|
||||||
@@ -757,6 +762,16 @@ impl ECStore {
|
|||||||
|
|
||||||
#[instrument(skip(self))]
|
#[instrument(skip(self))]
|
||||||
pub(super) async fn handle_delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
|
pub(super) async fn handle_delete_bucket(&self, bucket: &str, opts: &DeleteBucketOptions) -> Result<()> {
|
||||||
|
self.handle_delete_bucket_with_diagnostic_budget(bucket, opts, BucketDeleteDiagnosticBudget::new())
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn handle_delete_bucket_with_diagnostic_budget(
|
||||||
|
&self,
|
||||||
|
bucket: &str,
|
||||||
|
opts: &DeleteBucketOptions,
|
||||||
|
mut diagnostic_budget: BucketDeleteDiagnosticBudget,
|
||||||
|
) -> Result<()> {
|
||||||
if is_meta_bucketname(bucket) {
|
if is_meta_bucketname(bucket) {
|
||||||
return Err(StorageError::BucketNameInvalid(bucket.to_string()));
|
return Err(StorageError::BucketNameInvalid(bucket.to_string()));
|
||||||
}
|
}
|
||||||
@@ -815,14 +830,11 @@ impl ECStore {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
let mut diagnostic_budget = None;
|
|
||||||
|
|
||||||
if bucket_exists {
|
if bucket_exists {
|
||||||
validate_table_bucket_delete_guard(&self.ctx, bucket).await?;
|
validate_table_bucket_delete_guard(&self.ctx, bucket).await?;
|
||||||
|
|
||||||
if !opts.force {
|
if !opts.force {
|
||||||
let budget = diagnostic_budget.get_or_insert_with(BucketDeleteDiagnosticBudget::new);
|
if let Some(blocker) = bucket_delete_local_blocker(&self.ctx, bucket, &mut diagnostic_budget).await? {
|
||||||
if let Some(blocker) = bucket_delete_local_blocker(&self.ctx, bucket, budget).await? {
|
|
||||||
return Err(blocker);
|
return Err(blocker);
|
||||||
}
|
}
|
||||||
delete_opts.force_if_empty = true;
|
delete_opts.force_if_empty = true;
|
||||||
@@ -859,16 +871,11 @@ impl ECStore {
|
|||||||
if let Err(err) = delete_result
|
if let Err(err) = delete_result
|
||||||
&& (!sr_delete || !is_err_strict_volume_not_found(&err))
|
&& (!sr_delete || !is_err_strict_volume_not_found(&err))
|
||||||
{
|
{
|
||||||
if delete_opts.force_if_empty
|
if delete_opts.force_if_empty && matches!(&err, StorageError::BucketNotEmpty(_)) {
|
||||||
&& matches!(&err, StorageError::BucketNotEmpty(_))
|
let mut diagnostic_budget = BucketDeleteDiagnosticBudget::new();
|
||||||
&& let Some(blocker) = bucket_delete_local_blocker(
|
if let Some(blocker) = bucket_delete_local_blocker(&self.ctx, bucket, &mut diagnostic_budget).await? {
|
||||||
&self.ctx,
|
return Err(blocker);
|
||||||
bucket,
|
}
|
||||||
diagnostic_budget.get_or_insert_with(BucketDeleteDiagnosticBudget::new),
|
|
||||||
)
|
|
||||||
.await?
|
|
||||||
{
|
|
||||||
return Err(blocker);
|
|
||||||
}
|
}
|
||||||
return Err(err);
|
return Err(err);
|
||||||
}
|
}
|
||||||
@@ -895,10 +902,10 @@ impl ECStore {
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
BUCKET_DELETE_DIAGNOSTIC_MAX_ENTRIES, BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES, BucketDeleteBlockerKind,
|
BUCKET_DELETE_DIAGNOSTIC_MAX_ELAPSED, BUCKET_DELETE_DIAGNOSTIC_MAX_ENTRIES, BUCKET_DELETE_XLMETA_DIAGNOSTIC_MAX_BYTES,
|
||||||
BucketDeleteDiagnosticBudget, SCANNER_BUCKET_LIST_SET_CONCURRENCY, await_bucket_namespace_operation,
|
BucketDeleteBlockerKind, BucketDeleteDiagnosticBudget, SCANNER_BUCKET_LIST_SET_CONCURRENCY,
|
||||||
bucket_delete_metadata_cleanup_prefixes, bucket_deleted_marker_prefix, bucket_deleted_marker_volume,
|
await_bucket_namespace_operation, bucket_delete_metadata_cleanup_prefixes, bucket_deleted_marker_prefix,
|
||||||
run_bucket_usage_cleanup, run_physical_bucket_deletion, scan_metadata_less_residue,
|
bucket_deleted_marker_volume, run_bucket_usage_cleanup, run_physical_bucket_deletion, scan_metadata_less_residue,
|
||||||
scan_metadata_less_residue_with_budget, scanner_bucket_list_set_concurrency, should_override_created_from_metadata,
|
scan_metadata_less_residue_with_budget, scanner_bucket_list_set_concurrency, should_override_created_from_metadata,
|
||||||
validate_table_bucket_delete_allowed,
|
validate_table_bucket_delete_allowed,
|
||||||
};
|
};
|
||||||
@@ -2032,6 +2039,86 @@ mod tests {
|
|||||||
.expect("DeleteBucket must succeed once the client has drained the bucket");
|
.expect("DeleteBucket must succeed once the client has drained the bucket");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
#[serial]
|
||||||
|
async fn bucket_delete_defers_zero_evidence_diagnostic_timeout_to_physical_empty_check() {
|
||||||
|
let (disk_paths, ecstore) = setup_bucket_delete_test_env().await;
|
||||||
|
let bucket = format!("bucket-del-diag-{}", Uuid::new_v4().simple());
|
||||||
|
|
||||||
|
ecstore
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("empty bucket should be created");
|
||||||
|
let first_io_started = Arc::new(AtomicBool::new(false));
|
||||||
|
let diagnostic_budget = BucketDeleteDiagnosticBudget::new().with_first_io_delay(
|
||||||
|
BUCKET_DELETE_DIAGNOSTIC_MAX_ELAPSED + Duration::from_millis(100),
|
||||||
|
first_io_started.clone(),
|
||||||
|
);
|
||||||
|
|
||||||
|
ecstore
|
||||||
|
.handle_delete_bucket_with_diagnostic_budget(&bucket, &DeleteBucketOptions::default(), diagnostic_budget)
|
||||||
|
.await
|
||||||
|
.expect("a zero-evidence diagnostic timeout must defer to the physical empty check");
|
||||||
|
|
||||||
|
assert!(
|
||||||
|
first_io_started.load(Ordering::SeqCst),
|
||||||
|
"the regression must delay the first diagnostic read_dir"
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
!any_disk_path_exists(&disk_paths, &bucket).await,
|
||||||
|
"the physical empty check should allow deletion of the actually empty bucket"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
|
#[serial]
|
||||||
|
async fn bucket_delete_preserves_unobserved_residue_after_diagnostic_timeout() {
|
||||||
|
let (disk_paths, ecstore) = setup_bucket_delete_test_env().await;
|
||||||
|
let bucket = format!("bucket-del-residue-{}", Uuid::new_v4().simple());
|
||||||
|
let object = "object.txt";
|
||||||
|
|
||||||
|
ecstore
|
||||||
|
.make_bucket(&bucket, &MakeBucketOptions::default())
|
||||||
|
.await
|
||||||
|
.expect("bucket should be created");
|
||||||
|
write_metadata_less_part_on_all_disks(&disk_paths, &bucket, object).await;
|
||||||
|
let first_io_started = Arc::new(AtomicBool::new(false));
|
||||||
|
let diagnostic_budget = BucketDeleteDiagnosticBudget::new().with_first_io_delay(
|
||||||
|
BUCKET_DELETE_DIAGNOSTIC_MAX_ELAPSED + Duration::from_millis(100),
|
||||||
|
first_io_started.clone(),
|
||||||
|
);
|
||||||
|
|
||||||
|
let err = ecstore
|
||||||
|
.handle_delete_bucket_with_diagnostic_budget(&bucket, &DeleteBucketOptions::default(), diagnostic_budget)
|
||||||
|
.await
|
||||||
|
.expect_err("physical empty-bucket enforcement must reject unobserved residue");
|
||||||
|
assert!(
|
||||||
|
first_io_started.load(Ordering::SeqCst),
|
||||||
|
"the regression must time out before observing the residue"
|
||||||
|
);
|
||||||
|
match &err {
|
||||||
|
StorageError::BucketNotEmpty(err_bucket) | StorageError::BucketNotEmptyWithDetails { bucket: err_bucket, .. } => {
|
||||||
|
assert_eq!(err_bucket, &bucket)
|
||||||
|
}
|
||||||
|
other => panic!("expected an S3-compatible BucketNotEmpty error, got {other:?}"),
|
||||||
|
}
|
||||||
|
assert!(
|
||||||
|
any_disk_path_exists(&disk_paths, Path::new(&bucket).join(object)).await,
|
||||||
|
"physical empty-bucket enforcement must preserve unobserved residue"
|
||||||
|
);
|
||||||
|
|
||||||
|
ecstore
|
||||||
|
.delete_bucket(
|
||||||
|
&bucket,
|
||||||
|
&DeleteBucketOptions {
|
||||||
|
force: true,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("explicit force should clean up the test-owned residue");
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test(flavor = "multi_thread")]
|
#[tokio::test(flavor = "multi_thread")]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn bucket_delete_succeeds_after_listing_and_deleting_an_unversioned_overwrite() {
|
async fn bucket_delete_succeeds_after_listing_and_deleting_an_unversioned_overwrite() {
|
||||||
|
|||||||
@@ -97,6 +97,8 @@ pub(crate) struct BucketDeleteDiagnosticBudget {
|
|||||||
deadline: Option<tokio::time::Instant>,
|
deadline: Option<tokio::time::Instant>,
|
||||||
max_elapsed: Duration,
|
max_elapsed: Duration,
|
||||||
entries_remaining: usize,
|
entries_remaining: usize,
|
||||||
|
#[cfg(test)]
|
||||||
|
first_io_delay: Option<(Duration, Arc<std::sync::atomic::AtomicBool>)>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl BucketDeleteDiagnosticBudget {
|
impl BucketDeleteDiagnosticBudget {
|
||||||
@@ -109,9 +111,17 @@ impl BucketDeleteDiagnosticBudget {
|
|||||||
deadline: None,
|
deadline: None,
|
||||||
max_elapsed: elapsed,
|
max_elapsed: elapsed,
|
||||||
entries_remaining: entries,
|
entries_remaining: entries,
|
||||||
|
#[cfg(test)]
|
||||||
|
first_io_delay: None,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
fn with_first_io_delay(mut self, delay: Duration, started: Arc<std::sync::atomic::AtomicBool>) -> Self {
|
||||||
|
self.first_io_delay = Some((delay, started));
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
fn deadline(&mut self) -> tokio::time::Instant {
|
fn deadline(&mut self) -> tokio::time::Instant {
|
||||||
let max_elapsed = self.max_elapsed;
|
let max_elapsed = self.max_elapsed;
|
||||||
*self.deadline.get_or_insert_with(|| tokio::time::Instant::now() + max_elapsed)
|
*self.deadline.get_or_insert_with(|| tokio::time::Instant::now() + max_elapsed)
|
||||||
@@ -137,7 +147,20 @@ impl BucketDeleteDiagnosticBudget {
|
|||||||
if tokio::time::Instant::now() >= deadline {
|
if tokio::time::Instant::now() >= deadline {
|
||||||
return Ok(None);
|
return Ok(None);
|
||||||
}
|
}
|
||||||
match tokio::time::timeout_at(deadline, future).await {
|
#[cfg(test)]
|
||||||
|
let first_io_delay = self.first_io_delay.take();
|
||||||
|
#[cfg(test)]
|
||||||
|
let timeout_result = tokio::time::timeout_at(deadline, async move {
|
||||||
|
if let Some((delay, started)) = first_io_delay {
|
||||||
|
started.store(true, std::sync::atomic::Ordering::SeqCst);
|
||||||
|
tokio::time::sleep(delay).await;
|
||||||
|
}
|
||||||
|
future.await
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
#[cfg(not(test))]
|
||||||
|
let timeout_result = tokio::time::timeout_at(deadline, future).await;
|
||||||
|
match timeout_result {
|
||||||
Ok(result) => result.map(Some),
|
Ok(result) => result.map(Some),
|
||||||
Err(_) => Ok(None),
|
Err(_) => Ok(None),
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user