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 b4a9120173 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.
This commit is contained in:
Hauser
2026-09-30 19:24:32 +08:00
committed by GitHub
parent e573fb087b
commit eb350ad20d
4 changed files with 298 additions and 70 deletions
@@ -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::<Vec<_>>();
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
+25
View File
@@ -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]
+194 -69
View File
@@ -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 `<marker>` and `<marker>[` 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"<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>".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};
+29
View File
@@ -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<Self>,
bucket: &str,