From eb350ad20de5ca4f88a13c075e8db1e50c3caa22 Mon Sep 17 00:00:00 2001 From: Hauser Date: Wed, 30 Sep 2026 19:24:32 +0800 Subject: [PATCH] fix(ecstore): purge empty recursive bucket orphans (#8261) ## Related Issues Related to rustfs/backlog#2688. ## Summary of Changes Reclaim metadata-less orphan directory trees after a complete, first-page, empty recursive root listing in an authoritative never-versioned bucket. Both listing implementations use the same fail-closed cleanup decision. ## Verification Head b4a9120173d839dbd870e43e81d977a540a62c49 received an approval from loverustfs in review 5365460152. This merge does not add a local runtime validation claim. Full main CI remains a separate publication gate. ## Impact Empty recursive root listings can remove otherwise unaddressable physical residue. Versioned buckets and populated listings remain excluded; unreadable disks, object metadata, unknown files, and uncommitted data prevent cleanup. ## Additional Notes Reverting this change restores the previous orphan-directory behavior. --- .../src/set_disk/core/io_primitives.rs | 51 +++- crates/ecstore/src/set_disk/mod.rs | 25 ++ crates/ecstore/src/store/list_objects.rs | 263 +++++++++++++----- crates/ecstore/src/store/object.rs | 29 ++ 4 files changed, 298 insertions(+), 70 deletions(-) diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index b4e8db506..e73ee7332 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -6525,7 +6525,15 @@ impl SetDisks { let disks = self.get_disks_internal().await; - // Phase 1: classify every online disk. A directory that holds object + // An unresolved slot may hold the only copy of readable object + // metadata or uncommitted data under this prefix. Do not classify or + // remove orphan residue from the online subset while any drive is + // unavailable. + if disks.iter().any(Option::is_none) { + return Ok(false); + } + + // Phase 1: classify every disk. A directory that holds object // data or uncommitted residue on ANY disk blocks itself and its // ancestors on every disk, so a degraded/healable object is never // destroyed; purgeable subtrees beside it are still reclaimed. @@ -6579,6 +6587,47 @@ impl SetDisks { Ok(purged) } + /// Reclaim metadata-less directory trees discovered by a recursive empty + /// bucket listing. The initial directory read fails closed if any slot is + /// unavailable; per-prefix purges then reuse the full cross-disk scan. + pub(crate) async fn purge_orphan_dir_objects_in_bucket(&self, bucket: &str) -> bool { + let disks = self.get_disks_internal().await; + if disks.is_empty() || disks.iter().any(Option::is_none) { + return false; + } + + let mut prefixes = HashSet::new(); + for disk in disks.iter().flatten() { + let entries = match disk.list_dir("", bucket, "", 0).await { + Ok(entries) => entries, + Err(_) => return false, + }; + prefixes.extend( + entries + .into_iter() + .filter(|entry| entry.ends_with(SLASH_SEPARATOR) && is_safe_orphan_dir_entry(entry)), + ); + } + + let mut prefixes = prefixes.into_iter().collect::>(); + prefixes.sort_unstable(); + let mut purged = false; + for prefix in prefixes { + match self.purge_orphan_dir_object(bucket, &prefix).await { + Ok(prefix_purged) => purged |= prefix_purged, + Err(_) => return false, + } + } + purged + } + + /// Orphan cleanup must not infer a complete scan from only the online + /// subset of a set's disks. + pub(crate) async fn orphan_purge_has_complete_disk_set(&self) -> bool { + let disks = self.get_disks_internal().await; + !disks.is_empty() && disks.iter().all(Option::is_some) + } + fn orphan_purge_in_backoff(&self, key: &str) -> bool { let backoff = self .orphan_purge_backoff diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index ac2892bff..2b041ef13 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -9437,6 +9437,31 @@ mod tests { assert!(!purged, "a missing prefix should report nothing to purge"); } + #[tokio::test] + async fn orphan_directory_purge_preserves_tree_when_a_disk_slot_is_offline() { + let (dir, disk) = make_single_local_disk().await; + let prefix_dir = dir.path().join("bucket").join("pfx"); + fs::create_dir_all(prefix_dir.join("nested").join("leaf")) + .await + .expect("orphan directory tree should be created"); + + let set = make_set_disks_with(vec![Some(disk), None]).await; + let purged = set + .purge_orphan_dir_object("bucket", "pfx/") + .await + .expect("an unavailable slot should fail closed without a scan error"); + + assert!(!purged, "an incomplete disk scan must not claim the tree is an orphan"); + assert!(prefix_dir.join("nested/leaf").exists(), "online disk contents must remain untouched"); + + let bucket_purged = set.purge_orphan_dir_objects_in_bucket("bucket").await; + assert!(!bucket_purged, "bucket-wide cleanup must fail closed with an unavailable disk slot"); + assert!( + prefix_dir.join("nested/leaf").exists(), + "bucket-wide cleanup must preserve the online tree" + ); + } + // Cross-disk safety: if any drive still holds object data under the prefix, refuse // to purge on every drive so a degraded/healable object is never destroyed. #[tokio::test] diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index fc9121c2f..19fc2260d 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -387,6 +387,24 @@ fn should_purge_empty_directory_listing( && result.prefixes.is_empty() } +fn should_purge_empty_recursive_bucket_listing( + prefix: &str, + delimiter: Option<&str>, + marker: Option<&str>, + max_keys: i32, + incl_deleted: bool, + result: &ListObjectsInfo, +) -> bool { + prefix.is_empty() + && delimiter.is_none() + && marker.is_none() + && max_keys > 0 + && !incl_deleted + && !result.is_truncated + && result.objects.is_empty() + && result.prefixes.is_empty() +} + const MARKER_TAG_VERSION: &str = "v2"; const LEGACY_MARKER_TAG_VERSIONS: &[&str] = &["v1", MARKER_TAG_VERSION]; const LIST_CACHE_MARKER_PREFIX: &str = "[rustfs_cache:"; @@ -4006,82 +4024,91 @@ impl ECStore { // key sorting between `` and `[` is silently skipped on // the continuation page (backlog#1047). opts.parse_marker(); - if let Some(mode) = list_objects_index_mode_from_env() - && let Some(result) = self - .clone() + let key_only_result = if let Some(mode) = list_objects_index_mode_from_env() { + self.clone() .list_objects_from_opt_in_key_only_provider(&opts, mode, max_keys, incl_deleted) .await? - { - if should_purge_empty_directory_listing(prefix, opts.marker.as_deref(), max_keys, incl_deleted, &result) - && has_authoritative_never_versioned_state_in(&self.ctx, bucket) - .await - .unwrap_or(false) - { - self.purge_orphan_dir_object(bucket, prefix).await; - } - return Ok(result); - } - - let mut list_result = self - .clone() - .list_path(&opts) - .await - .unwrap_or_else(|err| MetaCacheEntriesSortedResult { - err: Some(to_filemeta_err(err)), - ..Default::default() - }); - let next_cache_id = list_result.entries.as_ref().and_then(|entries| entries.list_id.clone()); - - // err=None means gather_results filled its limit → disk has more data - let disk_has_more = list_result.err.is_none(); - - if let Some(err) = list_result.err.take() - && err != rustfs_filemeta::Error::Unexpected - { - return Err(to_object_err(err.into(), vec![bucket, prefix])); - } - - if let Some(result) = list_result.entries.as_mut() { - result.forward_past(opts.marker.clone()); - } - - // contextCanceled - - // Last RAW scanned key, captured before folding, so `list_objects_paginate` - // can advance past a fully-collapsed common-prefix page (ECA-03 / #944). - let last_scanned_key = last_scanned_entry_name(list_result.entries.as_ref()); - - let get_objects = ObjectInfo::from_meta_cache_entries_sorted_infos( - &list_result.entries.unwrap_or_default(), - bucket, - prefix, - delimiter.clone(), - ) - .await; - - let (objects, prefixes, is_truncated, next_marker, next_version_idmarker) = list_objects_paginate( - get_objects, - &delimiter, - max_keys, - disk_has_more, - next_cache_id.as_deref(), - false, - last_scanned_key.as_deref(), - ); - let _ = next_version_idmarker; - - let result = ListObjectsInfo { - is_truncated, - next_marker, - objects, - prefixes, + } else { + None }; - if should_purge_empty_directory_listing(prefix, opts.marker.as_deref(), max_keys, incl_deleted, &result) + + let result = if let Some(result) = key_only_result { + result + } else { + let mut list_result = self + .clone() + .list_path(&opts) + .await + .unwrap_or_else(|err| MetaCacheEntriesSortedResult { + err: Some(to_filemeta_err(err)), + ..Default::default() + }); + let next_cache_id = list_result.entries.as_ref().and_then(|entries| entries.list_id.clone()); + + // err=None means gather_results filled its limit → disk has more data + let disk_has_more = list_result.err.is_none(); + + if let Some(err) = list_result.err.take() + && err != rustfs_filemeta::Error::Unexpected + { + return Err(to_object_err(err.into(), vec![bucket, prefix])); + } + + if let Some(result) = list_result.entries.as_mut() { + result.forward_past(opts.marker.clone()); + } + + // Last RAW scanned key, captured before folding, so `list_objects_paginate` + // can advance past a fully-collapsed common-prefix page (ECA-03 / #944). + let last_scanned_key = last_scanned_entry_name(list_result.entries.as_ref()); + + let get_objects = ObjectInfo::from_meta_cache_entries_sorted_infos( + &list_result.entries.unwrap_or_default(), + bucket, + prefix, + delimiter.clone(), + ) + .await; + + let (objects, prefixes, is_truncated, next_marker, next_version_idmarker) = list_objects_paginate( + get_objects, + &delimiter, + max_keys, + disk_has_more, + next_cache_id.as_deref(), + false, + last_scanned_key.as_deref(), + ); + let _ = next_version_idmarker; + + ListObjectsInfo { + is_truncated, + next_marker, + objects, + prefixes, + } + }; + + let purge_exact_prefix = + should_purge_empty_directory_listing(prefix, opts.marker.as_deref(), max_keys, incl_deleted, &result); + let purge_empty_bucket = should_purge_empty_recursive_bucket_listing( + prefix, + delimiter.as_deref(), + opts.marker.as_deref(), + max_keys, + incl_deleted, + &result, + ); + if (purge_exact_prefix || purge_empty_bucket) && has_authoritative_never_versioned_state_in(&self.ctx, bucket) .await .unwrap_or(false) { - self.purge_orphan_dir_object(bucket, prefix).await; + if purge_exact_prefix { + self.purge_orphan_dir_object(bucket, prefix).await; + } else { + self.purge_orphan_dir_objects_in_bucket(bucket).await; + } } Ok(result) } @@ -9507,6 +9534,38 @@ mod test { assert!(!should_purge_empty_directory_listing("ghost/", None, 1, false, &truncated)); } + #[test] + fn recursive_bucket_orphan_purge_requires_an_empty_complete_root_scan() { + let empty = ListObjectsInfo::default(); + assert!(super::should_purge_empty_recursive_bucket_listing("", None, None, 1, false, &empty)); + assert!(!super::should_purge_empty_recursive_bucket_listing( + "ghost/", None, None, 1, false, &empty + )); + assert!(!super::should_purge_empty_recursive_bucket_listing("", Some("/"), None, 1, false, &empty)); + assert!(!super::should_purge_empty_recursive_bucket_listing( + "", + None, + Some("marker"), + 1, + false, + &empty + )); + assert!(!super::should_purge_empty_recursive_bucket_listing("", None, None, 0, false, &empty)); + assert!(!super::should_purge_empty_recursive_bucket_listing("", None, None, 1, true, &empty)); + + let live = ListObjectsInfo { + objects: vec![ObjectInfo::default()], + ..Default::default() + }; + assert!(!super::should_purge_empty_recursive_bucket_listing("", None, None, 1, false, &live)); + + let incomplete = ListObjectsInfo { + is_truncated: true, + ..Default::default() + }; + assert!(!super::should_purge_empty_recursive_bucket_listing("", None, None, 1, false, &incomplete)); + } + #[tokio::test] async fn empty_recursive_listing_purges_committed_delete_residue() { use crate::bucket::metadata_sys::{init_bucket_metadata_sys, test_support::isolated_store_over_temp_disks}; @@ -9558,6 +9617,72 @@ mod test { } } + #[tokio::test] + async fn empty_recursive_bucket_listing_purges_orphan_directory_prefixes() { + use crate::bucket::metadata_sys::{init_bucket_metadata_sys, test_support::isolated_store_over_temp_disks}; + use crate::storage_api_contracts::bucket::{BucketOperations as _, MakeBucketOptions}; + + let (dirs, store) = isolated_store_over_temp_disks().await; + let bucket = "recursive-bucket-orphan-purge"; + init_bucket_metadata_sys(store.clone(), Vec::new()).await; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created with authoritative metadata"); + + for dir in &dirs { + tokio::fs::create_dir_all(dir.path().join(bucket).join("ghost").join("nested").join("leaf")) + .await + .expect("metadata-less orphan directory tree should be created"); + } + + let result = store + .clone() + .list_objects_generic(bucket, "", None, None, 1000, false) + .await + .expect("recursive bucket listing should succeed"); + + assert!(result.objects.is_empty(), "orphan directories are not S3 objects"); + assert!(result.prefixes.is_empty(), "a delimiter-less listing has no CommonPrefixes"); + for dir in &dirs { + assert!( + !dir.path().join(bucket).join("ghost").exists(), + "an empty recursive bucket scan should reclaim its metadata-less orphan prefix" + ); + } + + let versioned_bucket = "recursive-bucket-versioned-orphan"; + store + .make_bucket(versioned_bucket, &MakeBucketOptions::default()) + .await + .expect("versioned test bucket should be created"); + store + .update_bucket_metadata_config( + versioned_bucket, + crate::bucket::metadata::BUCKET_VERSIONING_CONFIG, + b"Enabled".to_vec(), + ) + .await + .expect("bucket versioning should be enabled"); + for dir in &dirs { + tokio::fs::create_dir_all(dir.path().join(versioned_bucket).join("ghost").join("nested").join("leaf")) + .await + .expect("versioned bucket orphan tree should be created"); + } + + store + .clone() + .list_objects_generic(versioned_bucket, "", None, None, 1000, false) + .await + .expect("recursive versioned bucket listing should succeed"); + for dir in &dirs { + assert!( + dir.path().join(versioned_bucket).join("ghost").exists(), + "root recursive LIST must not purge orphan residue in a versioned bucket" + ); + } + } + #[tokio::test] async fn list_objects_exact_prefix_paginates_across_storage_layers() { use crate::bucket::metadata_sys::{init_bucket_metadata_sys, test_support::isolated_store_over_temp_disks}; diff --git a/crates/ecstore/src/store/object.rs b/crates/ecstore/src/store/object.rs index f01e31d3b..e58b0b2d1 100644 --- a/crates/ecstore/src/store/object.rs +++ b/crates/ecstore/src/store/object.rs @@ -4693,6 +4693,14 @@ impl ECStore { /// the caller falls back to surfacing the original NotFound. pub(super) async fn purge_orphan_dir_object(&self, bucket: &str, object: &str) -> bool { let prefix = decode_dir_object(object); + for pool in self.pools.iter() { + for set in pool.disk_set.iter() { + if !set.orphan_purge_has_complete_disk_set().await { + return false; + } + } + } + let mut purged = false; for pool in self.pools.iter() { for set in pool.disk_set.iter() { @@ -4713,6 +4721,27 @@ impl ECStore { purged } + pub(super) async fn purge_orphan_dir_objects_in_bucket(&self, bucket: &str) -> bool { + if is_meta_bucketname(bucket) { + return false; + } + for pool in self.pools.iter() { + for set in pool.disk_set.iter() { + if !set.orphan_purge_has_complete_disk_set().await { + return false; + } + } + } + + let mut purged = false; + for pool in self.pools.iter() { + for set in pool.disk_set.iter() { + purged |= set.purge_orphan_dir_objects_in_bucket(bucket).await; + } + } + purged + } + pub async fn delete_object_with_tier_delete_journal( self: &Arc, bucket: &str,