diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 5b0c5a7ca..12099907b 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -94,6 +94,8 @@ const STALE_TMP_OBJECT_EXPIRY: Duration = Duration::from_secs(24 * 60 * 60); #[cfg(test)] tokio::task_local! { static DIRECTORY_LISTING_ENTRY_PROBE_COUNT: Arc; + /// `(sampled, complete)` directory reads issued by the delete-residue probe. + static DELETE_RESIDUE_PROBE_READS: Arc<(AtomicUsize, AtomicUsize)>; } const RUSTFS_META_TMP_OLD_BUCKET: &str = ".rustfs.sys/tmp-old"; const INLINE_METADATA_ROLLBACK_DIR_XOR: u128 = 0x7275737466735f696e6c696e655f7262; @@ -105,6 +107,21 @@ pub(crate) const RESERVED_DELETE_DATA_DIR_MARKER_PREFIX: &str = "reserve-delete- /// under-filled batch settles the common case without materializing large /// child sets; a full batch cannot prove no listable child hides behind it. const DELETE_RESIDUE_PROBE_LIMIT: i32 = 8; +/// Most directory reads one delete-residue probe spends walking down through +/// the ancestors of a deleted key before it gives up and lets the prefix +/// surface. A genuine prefix settles within its depth (the first object's +/// `xl.meta` ends the walk), so the budget only caps the cost of hiding a +/// large residue tree, which then hides one level at a time instead. +const DELETE_RESIDUE_PROBE_READ_BUDGET: usize = 32; + +/// One directory read owed by the delete-residue probe. +enum ResidueProbeStep { + /// Read a bounded batch of `dir` and probe what it holds. + Sample(String), + /// Read all of `dir` once every child in `sampled` proved unlistable, and + /// probe the rest. + Remainder { dir: String, sampled: Vec }, +} /// A `part.N` file with a positive part number, the shape erasure data takes /// inside a version data dir. @@ -8047,72 +8064,116 @@ impl LocalDisk { Ok(false) } - /// Whether the metadata-less directory `dir_name` holds nothing but the - /// data dirs of deleted versions: it is itself a non-nil UUID directory of - /// `part.N` files and delete-transaction markers, or every child is one. - /// That is what an interrupted or deferred version delete leaves behind - /// once the `xl.meta` is gone, and it must not surface as a prefix. Real - /// object children are directories carrying their own `xl.meta`, so the - /// first non-UUID child, stray file, or subdirectory inside a UUID child - /// proves the directory is a genuine prefix. Reads are bounded: a + /// Whether the metadata-less directory `dir_name` holds nothing that a + /// listing could show: every leaf under it is the data dir of a deleted + /// version (a non-nil UUID directory of `part.N` files and + /// delete-transaction markers) or an empty directory. That is what an + /// interrupted or deferred version delete leaves behind once the `xl.meta` + /// is gone, and neither the object directory nor the date-style ancestors + /// above it may surface as prefixes (#6898). Real objects are directories + /// carrying their own `xl.meta`, so the first file met outside a UUID data + /// dir, or any stray entry inside one, proves a genuine prefix and ends the + /// walk at the depth of the first object. + /// + /// Reads stay small on genuine prefixes: each directory is sampled with one + /// bounded batch, the sampled children are probed depth-first, and the + /// directory is only read in full once every sampled child proved to hold + /// nothing listable. A total read budget makes a residue tree larger than + /// the budget surface and be hidden one level down instead, so the cost on + /// a genuine prefix never exceeds its depth in small directory reads. A /// directory that vanishes mid-probe holds nothing listable. async fn directory_is_delete_residue(&self, bucket: &str, dir_name: &str, stall: Option) -> Result { let dir_name = dir_name.trim_end_matches(SLASH_SEPARATOR); - let Some(entries) = self.read_dir_for_residue_probe(bucket, dir_name, stall).await? else { - return Ok(false); - }; - if entries.is_empty() { - return Ok(false); - } - - let is_data_dir = dir_name - .rsplit(SLASH_SEPARATOR) - .next() - .is_some_and(|name| Uuid::parse_str(name).is_ok_and(|uuid| !uuid.is_nil())); - if is_data_dir && entries.iter().all(|entry| is_metadata_less_data_dir_entry(entry)) { - return Ok(true); - } - - for entry in entries { - let Some(child) = entry.strip_suffix(SLASH_SEPARATOR) else { - return Ok(false); - }; - if !Uuid::parse_str(child).is_ok_and(|uuid| !uuid.is_nil()) { + let mut reads = 0usize; + let mut pending = vec![ResidueProbeStep::Sample(dir_name.to_owned())]; + while let Some(step) = pending.pop() { + if reads >= DELETE_RESIDUE_PROBE_READ_BUDGET { return Ok(false); } + reads += 1; - let child_path = path_join_buf(&[dir_name, child]); - let Some(child_entries) = self.read_dir_for_residue_probe(bucket, &child_path, stall).await? else { + let (dir, entries, complete) = match step { + ResidueProbeStep::Sample(dir) => { + let Some((entries, complete)) = self.read_dir_for_residue_probe(bucket, &dir, stall, false).await? else { + continue; + }; + (dir, entries, complete) + } + ResidueProbeStep::Remainder { dir, sampled } => { + let Some((entries, _)) = self.read_dir_for_residue_probe(bucket, &dir, stall, true).await? else { + continue; + }; + let entries = entries + .into_iter() + .filter(|entry| !sampled.contains(entry)) + .collect::>(); + (dir, entries, true) + } + }; + + let is_data_dir = dir + .rsplit(SLASH_SEPARATOR) + .next() + .is_some_and(|name| Uuid::parse_str(name).is_ok_and(|uuid| !uuid.is_nil())); + if is_data_dir { + if !entries.iter().all(|entry| is_metadata_less_data_dir_entry(entry)) { + // A subdirectory, an `xl.meta`, or an unknown file inside a + // UUID directory: not plain delete residue. + return Ok(false); + } + if !complete { + // Only the sampled part files were seen; the rest must be + // read before the data dir counts as plain residue. + pending.push(ResidueProbeStep::Remainder { dir, sampled: entries }); + } continue; - }; - if !child_entries.iter().all(|entry| is_metadata_less_data_dir_entry(entry)) { - return Ok(false); } + + let mut children = Vec::with_capacity(entries.len()); + for entry in &entries { + let Some(child) = entry.strip_suffix(SLASH_SEPARATOR) else { + // A file outside a plain data dir: `xl.meta` or something + // this probe does not understand. Either way, not residue. + return Ok(false); + }; + children.push(path_join_buf(&[&dir, child])); + } + if !complete { + // Revisit the unsampled siblings only after every sampled + // child, probed first, turned out to hold nothing listable. + pending.push(ResidueProbeStep::Remainder { dir, sampled: entries }); + } + pending.extend(children.into_iter().map(ResidueProbeStep::Sample)); } Ok(true) } - /// Read `dir` with a bounded batch first and a complete read only when the - /// batch was full. `None` when the directory does not exist any more. - async fn read_dir_for_residue_probe(&self, bucket: &str, dir: &str, stall: Option) -> Result>> { - for count in [DELETE_RESIDUE_PROBE_LIMIT, -1] { - let entries = match with_walk_stall_timeout(stall, self.list_dir("", bucket, dir, count)).await { - Ok(entries) => entries, - Err(err) => { - if err == DiskError::VolumeNotFound || err == Error::FileNotFound { - return Ok(None); - } - - return Err(err); - } - }; - if count < 0 || entries.len() < count as usize { - return Ok(Some(entries)); + /// Read `dir` for the residue probe: one bounded batch, or the complete + /// directory when `complete` is set. The flag in the result says whether + /// the batch held the whole directory. `None` when the directory does not + /// exist any more. + async fn read_dir_for_residue_probe( + &self, + bucket: &str, + dir: &str, + stall: Option, + complete: bool, + ) -> Result, bool)>> { + let count = if complete { -1 } else { DELETE_RESIDUE_PROBE_LIMIT }; + #[cfg(test)] + let _ = DELETE_RESIDUE_PROBE_READS.try_with(|reads| { + let counter = if complete { &reads.1 } else { &reads.0 }; + counter.fetch_add(1, Ordering::Relaxed) + }); + match with_walk_stall_timeout(stall, self.list_dir("", bucket, dir, count)).await { + Ok(entries) => { + let complete = complete || entries.len() < DELETE_RESIDUE_PROBE_LIMIT as usize; + Ok(Some((entries, complete))) } + Err(err) if err == DiskError::VolumeNotFound || err == Error::FileNotFound => Ok(None), + Err(err) => Err(err), } - - Ok(None) } /// Whether anything under `dir_name` would appear in a listing. With @@ -18719,9 +18780,9 @@ mod test { let (fast_path_names, fast_path_probes) = scan_prefixes(&disk, bucket, true).await; assert_eq!(conservative_names, expected_names); - let mut expected_fast_path_names = expected_names.clone(); - expected_fast_path_names.push("stale/".to_owned()); - assert_eq!(fast_path_names, expected_fast_path_names); + // The fast path hides the empty `stale/` chain too, at the cost of a + // few bounded directory reads rather than the metadata probes. + assert_eq!(fast_path_names, expected_names); let expected_probes = PREFIX_COUNT * 3 + 3; assert_eq!(conservative_probes, expected_probes); assert_eq!(fast_path_probes, 0); @@ -18832,9 +18893,9 @@ mod test { } // Directories whose only content is a deleted version's data dir are - // not prefixes; their ancestors stay ordinary directories until an - // empty listing reclaims them. + // not prefixes, and neither are their ancestors. assert_eq!(scan_names(&disk, bucket, "residue/2026/").await, Vec::::new()); + assert_eq!(scan_names(&disk, bucket, "residue/").await, Vec::::new()); assert_eq!(scan_names(&disk, bucket, "committed/").await, Vec::::new()); // UUID-named directories holding real objects, an object keyed by a @@ -18842,18 +18903,258 @@ mod test { assert_eq!(scan_names(&disk, bucket, "uploads/").await, vec![format!("uploads/{upload}/")]); assert_eq!(scan_names(&disk, bucket, "named/").await, vec![format!("named/{named}")]); assert_eq!(scan_names(&disk, bucket, "mixed/").await, vec!["mixed/child".to_owned()]); + assert_eq!( + scan_names(&disk, bucket, "").await, + vec!["mixed/".to_owned(), "named/".to_owned(), "uploads/".to_owned()] + ); + } + + #[tokio::test] + async fn test_scan_dir_nonrecursive_fast_path_hides_delete_residue_ancestors() { + use rustfs_filemeta::MetacacheReader; + use tempfile::tempdir; + + let dir = tempdir().expect("tempdir should be created"); + let bucket = "test-bucket"; + let bucket_dir = dir.path().join(bucket); + + async fn write_object(object_dir: &Path, object_name: &str) { + fs::create_dir_all(object_dir) + .await + .expect("object directory should be created"); + let mut metadata = FileMeta::default(); + let mut file_info = FileInfo::new(object_name, 1, 1); + file_info.mod_time = Some(OffsetDateTime::now_utc()); + metadata.add_version(file_info).expect("metadata should be valid"); + fs::write( + object_dir.join(STORAGE_FORMAT_FILE), + metadata.marshal_msg().expect("metadata should encode"), + ) + .await + .expect("object metadata should be written"); + } + + async fn write_residue(object_dir: &Path, committed: bool) { + let residue = object_dir.join(Uuid::new_v4().to_string()); + fs::create_dir_all(&residue).await.expect("residue should be created"); + fs::write(residue.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + if committed { + fs::write(residue.join(format!("{DELETE_DATA_DIR_MARKER_PREFIX}{}", Uuid::new_v4())), []) + .await + .expect("delete marker should be written"); + } + } + + // The reported shape: a date-partitioned key deleted on an older build, + // whose data dir survived without a delete-transaction marker. Every + // ancestor up to `metrics/` holds nothing else. + write_residue(&bucket_dir.join("metrics/kubelet/2026/08/28/23/74992556388248657933757.parquet"), false).await; + // Several committed residues under one hour plus an empty sibling hour. + write_residue(&bucket_dir.join("metrics/cpu/2026/08/28/22/a.parquet"), true).await; + write_residue(&bucket_dir.join("metrics/cpu/2026/08/28/22/b.parquet"), true).await; + fs::create_dir_all(bucket_dir.join("metrics/cpu/2026/08/28/21")) + .await + .expect("empty hour directory should be created"); + + // A live object deep under an otherwise identical tree keeps every + // ancestor visible, even beside residue. + write_residue(&bucket_dir.join("logs/default/2026/08/28/23/old.parquet"), false).await; + write_object( + &bucket_dir.join("logs/default/2026/08/28/23/live.parquet"), + "logs/default/2026/08/28/23/live.parquet", + ) + .await; + + // A prefix with more residue directories than the probe budget stays + // visible rather than costing an unbounded walk. + for hour in 0..(DELETE_RESIDUE_PROBE_READ_BUDGET + 1) { + write_residue(&bucket_dir.join(format!("bulk/2026/08/28/{hour:02}/a.parquet")), false).await; + } + // Exactly at the budget: `edge/` with N object dirs costs one sampled + // read of `edge`, two reads per object dir (its own and its data dir), + // and one remainder read of `edge` once its 8-entry sample was full. + for object in 0..15 { + write_residue(&bucket_dir.join(format!("edge-hide/{object:02}.parquet")), false).await; + } + for object in 0..16 { + write_residue(&bucket_dir.join(format!("edge-show/{object:02}.parquet")), false).await; + } + + // A file that is not `xl.meta` outside a UUID data dir is not residue + // either: the probe does not guess, the prefix surfaces. + fs::create_dir_all(bucket_dir.join("stray/2026/08")) + .await + .expect("stray directory should be created"); + fs::write(bucket_dir.join("stray/2026/08/notes.txt"), b"?") + .await + .expect("stray file should be written"); + + // A UUID directory holding a subdirectory is not a data dir shape the + // probe understands, so it surfaces even when the subdirectory is empty. + let odd = Uuid::new_v4().to_string(); + fs::create_dir_all(bucket_dir.join("odd").join(&odd).join("sub")) + .await + .expect("odd directory should be created"); + + let endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be UTF-8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should initialize"); + + async fn scan_names(disk: &LocalDisk, bucket: &str, current: &str) -> Vec { + let (reader, mut writer) = tokio::io::duplex(64 * 1024); + let mut output = MetacacheWriter::new(&mut writer); + let opts = WalkDirOptions { + bucket: bucket.to_string(), + base_dir: current.to_string(), + skip_hidden_prefix_check: true, + ..Default::default() + }; + let mut objects_returned = 0; + disk.scan_dir( + current.to_string(), + "".to_string(), + &opts, + &mut output, + &mut objects_returned, + false, + None, + ) + .await + .expect("scan_dir should succeed"); + output.close().await.expect("metacache writer should close"); + drop(output); + drop(writer); + + let mut names = MetacacheReader::new(reader) + .read_all() + .await + .expect("scan output should decode") + .into_iter() + .map(|entry| entry.name) + .collect::>(); + names.sort(); + names + } + + // Every level of a tree whose only leaves are deleted data dirs is + // hidden, not just the object directory itself. + assert_eq!(scan_names(&disk, bucket, "metrics/kubelet/2026/08/28/").await, Vec::::new()); + assert_eq!(scan_names(&disk, bucket, "metrics/kubelet/").await, Vec::::new()); + assert_eq!(scan_names(&disk, bucket, "metrics/cpu/2026/08/28/").await, Vec::::new()); + assert_eq!(scan_names(&disk, bucket, "metrics/").await, Vec::::new()); + assert_eq!( + scan_names(&disk, bucket, "logs/default/2026/08/28/").await, + vec!["logs/default/2026/08/28/23/".to_owned()] + ); + assert_eq!(scan_names(&disk, bucket, "logs/").await, vec!["logs/default/".to_owned()]); + assert_eq!(scan_names(&disk, bucket, "bulk/2026/08/").await, vec!["bulk/2026/08/28/".to_owned()]); + // Listed directly, every object dir is probed on its own and hidden; + // the budget only bites where the whole tree hangs off one entry, so + // the root listing below shows `edge-show/` but not `edge-hide/`. + assert_eq!(scan_names(&disk, bucket, "edge-hide/").await, Vec::::new()); + assert_eq!(scan_names(&disk, bucket, "edge-show/").await, Vec::::new()); + assert_eq!(scan_names(&disk, bucket, "stray/").await, vec!["stray/2026/".to_owned()]); + assert_eq!(scan_names(&disk, bucket, "odd/").await, vec![format!("odd/{odd}/")]); assert_eq!( scan_names(&disk, bucket, "").await, vec![ - "committed/".to_owned(), - "mixed/".to_owned(), - "named/".to_owned(), - "residue/".to_owned(), - "uploads/".to_owned(), + "bulk/".to_owned(), + "edge-show/".to_owned(), + "logs/".to_owned(), + "odd/".to_owned(), + "stray/".to_owned(), ] ); } + #[tokio::test] + async fn test_scan_dir_nonrecursive_fast_path_residue_probe_reads_stay_within_prefix_depth() { + use rustfs_filemeta::MetacacheReader; + use tempfile::tempdir; + + let dir = tempdir().expect("tempdir should be created"); + let bucket = "test-bucket"; + let bucket_dir = dir.path().join(bucket); + + async fn write_object(object_dir: &Path, object_name: &str) { + fs::create_dir_all(object_dir) + .await + .expect("object directory should be created"); + let mut metadata = FileMeta::default(); + let mut file_info = FileInfo::new(object_name, 1, 1); + file_info.mod_time = Some(OffsetDateTime::now_utc()); + metadata.add_version(file_info).expect("metadata should be valid"); + fs::write( + object_dir.join(STORAGE_FORMAT_FILE), + metadata.marshal_msg().expect("metadata should encode"), + ) + .await + .expect("object metadata should be written"); + } + + // A genuine date-partitioned prefix wider than the probe batch at two + // levels: 20 days, each hour holding 100 objects. + const DAYS: usize = 20; + const OBJECTS: usize = 100; + for day in 0..DAYS { + for object in 0..OBJECTS { + let key = format!("data/stream/2026/08/{day:02}/23/{object:03}.parquet"); + write_object(&bucket_dir.join(&key), &key).await; + } + } + + let endpoint = + Endpoint::try_from(dir.path().to_str().expect("tempdir path should be UTF-8")).expect("endpoint should parse"); + let disk = LocalDisk::new(&endpoint, false).await.expect("local disk should initialize"); + + let reads = Arc::new((AtomicUsize::new(0), AtomicUsize::new(0))); + let (reader, mut writer) = tokio::io::duplex(64 * 1024); + let mut output = MetacacheWriter::new(&mut writer); + let opts = WalkDirOptions { + bucket: bucket.to_string(), + base_dir: "data/".to_string(), + skip_hidden_prefix_check: true, + ..Default::default() + }; + let mut objects_returned = 0; + DELETE_RESIDUE_PROBE_READS + .scope( + Arc::clone(&reads), + disk.scan_dir( + "data/".to_string(), + "".to_string(), + &opts, + &mut output, + &mut objects_returned, + false, + None, + ), + ) + .await + .expect("scan_dir should succeed"); + output.close().await.expect("metacache writer should close"); + drop(output); + drop(writer); + + let names = MetacacheReader::new(reader) + .read_all() + .await + .expect("scan output should decode") + .into_iter() + .map(|entry| entry.name) + .collect::>(); + assert_eq!(names, vec!["data/stream/".to_owned()]); + + // stream, 2026, 08, one sampled day, its hour, one sampled object: + // the walk ends at the first `xl.meta`, one bounded read per level, + // without ever reading a wide directory in full. + let (sampled, complete) = (reads.0.load(Ordering::Relaxed), reads.1.load(Ordering::Relaxed)); + assert_eq!(sampled, 6, "one sampled read per level down to the first object"); + assert_eq!(complete, 0, "a genuine prefix must never cost a complete directory read"); + } + #[tokio::test] async fn test_scan_dir_nonrecursive_skips_dirs_with_only_hidden_delete_markers() { use rustfs_filemeta::MetacacheReader; diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index e2cef14a2..ae8694826 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -3494,19 +3494,44 @@ fn dangling_delete_grace() -> time::Duration { /// Result of scanning one disk's copy of a directory prefix while deciding /// whether an orphan (metadata-less) directory tree can be safely purged. enum OrphanDirScan { - /// The subtree holds object metadata or uncommitted data, so it must not be - /// purged. - HasData, - /// The prefix contains only empty directories and/or UUID data directories - /// carrying a committed delete marker. - Purgeable { - empty_dirs: Vec, + /// The prefix exists on this disk. `blocked` holds every directory that is + /// itself, or an ancestor of, object metadata or uncommitted data (closed + /// under taking parents, root included when anything under it is blocked); + /// `dirs` is the pre-order list of every directory reached that is not + /// blocked, and `committed_files` the erasure data and committed delete + /// markers found in the unblocked UUID data dirs among them. + Scanned { + blocked: HashSet, + dirs: Vec, committed_files: Vec, }, + /// A directory read failed for a reason other than absence, so nothing + /// under the prefix can be classified on this disk. + Unreadable, /// The prefix does not exist on this disk. Missing, } +/// How long an orphan prefix whose purge scan met unpurgeable data is left +/// alone before an empty listing scans it again. +const ORPHAN_PURGE_BACKOFF: Duration = Duration::from_secs(60); +/// Upper bound on remembered backoff entries; the oldest is dropped first. +const ORPHAN_PURGE_BACKOFF_MAX_ENTRIES: usize = 4096; + +/// Mark `dir` and every ancestor up to and including `root` as blocked. +fn block_orphan_dir_chain(blocked: &mut HashSet, root: &str, dir: &str) { + let mut current = dir; + loop { + if !blocked.insert(current.to_owned()) || current == root { + return; + } + let Some((parent, _)) = current.rsplit_once(SLASH_SEPARATOR) else { + return; + }; + current = parent; + } +} + fn is_safe_orphan_dir_entry(entry: &str) -> bool { let component = entry.strip_suffix(SLASH_SEPARATOR).unwrap_or(entry); !component.is_empty() @@ -6311,9 +6336,12 @@ impl SetDisks { .await } - /// Scan a single disk's copy of `prefix` and decide whether it is an orphan - /// directory subtree. Only empty directories and UUID data directories with - /// valid committed delete markers are purgeable; every child is still scanned. + /// Scan a single disk's copy of `prefix` and classify every directory + /// under it. Only empty directories and UUID data directories with valid + /// committed delete markers are purgeable; anything else blocks its whole + /// ancestor chain while sibling subtrees stay purgeable, so committed + /// residue is still reclaimed when it shares an ancestor with residue from + /// an older build that never wrote markers (#6898). async fn scan_orphan_dir(disk: &DiskStore, bucket: &str, prefix: &str) -> OrphanDirScan { let root = prefix.trim_end_matches(SLASH_SEPARATOR).to_string(); let mut stack = vec![root.clone()]; @@ -6321,6 +6349,7 @@ impl SetDisks { // so reversing it yields a safe children-first removal order. let mut dirs: Vec = Vec::new(); let mut committed_files: Vec = Vec::new(); + let mut blocked: HashSet = HashSet::new(); let mut existed = false; while let Some(dir) = stack.pop() { @@ -6335,16 +6364,18 @@ impl SetDisks { } // Classification must fail closed: committed residue is safe to // remove only after every reachable child was inspected. - Err(_) => return OrphanDirScan::HasData, + Err(_) => return OrphanDirScan::Unreadable, }; existed = true; let mut child_dirs = Vec::new(); let mut files = Vec::new(); + let mut has_data = false; for entry in entries { if !is_safe_orphan_dir_entry(&entry) { - return OrphanDirScan::HasData; + has_data = true; + break; } match entry.strip_suffix(SLASH_SEPARATOR) { Some(child) => child_dirs.push(format!("{dir}{SLASH_SEPARATOR}{child}")), @@ -6352,28 +6383,29 @@ impl SetDisks { } } - if !files.is_empty() { + if !has_data && !files.is_empty() { let data_dir_name = dir.rsplit(SLASH_SEPARATOR).next().unwrap_or_default(); let is_uuid_data_dir = Uuid::parse_str(data_dir_name).is_ok_and(|uuid| !uuid.is_nil()); let has_committed_delete = files.iter().any(|entry| is_committed_delete_marker(entry)); + has_data = !is_uuid_data_dir || !has_committed_delete || files.iter().any(|entry| entry == STORAGE_FORMAT_FILE); + } - if !is_uuid_data_dir || !has_committed_delete || files.iter().any(|entry| entry == STORAGE_FORMAT_FILE) { - return OrphanDirScan::HasData; - } - - committed_files.extend(files.into_iter().map(|entry| path_join_buf(&[&dir, &entry]))); - dirs.push(dir); - stack.extend(child_dirs); + if has_data { + // Nothing below a blocked directory is ever removed, so its + // children need no classification. + block_orphan_dir_chain(&mut blocked, &root, &dir); continue; } + committed_files.extend(files.into_iter().map(|entry| format!("{dir}{SLASH_SEPARATOR}{entry}"))); dirs.push(dir); stack.extend(child_dirs); } if existed { - OrphanDirScan::Purgeable { - empty_dirs: dirs, + OrphanDirScan::Scanned { + blocked, + dirs, committed_files, } } else { @@ -6468,43 +6500,101 @@ impl SetDisks { /// of this set (the caller should surface the original NotFound), and `Err` on /// a hard disk failure. pub(crate) async fn purge_orphan_dir_object(&self, bucket: &str, object: &str) -> disk::error::Result { + let backoff_key = format!("{bucket}{SLASH_SEPARATOR}{object}"); + if self.orphan_purge_in_backoff(&backoff_key) { + return Ok(false); + } + let disks = self.get_disks_internal().await; - // Phase 1: classify every online disk. Refuse to purge if ANY disk holds - // object data under the prefix, so a degraded/healable object is never - // destroyed. - let mut per_disk_dirs: Vec<(usize, Vec, Vec)> = Vec::new(); - let mut existed = false; + // Phase 1: classify every online 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. + let mut per_disk: Vec<(usize, Vec, Vec)> = Vec::new(); + let mut blocked: HashSet = HashSet::new(); for (i, disk) in disks.iter().enumerate() { let Some(disk) = disk else { continue }; match Self::scan_orphan_dir(disk, bucket, object).await { - OrphanDirScan::HasData => return Ok(false), - OrphanDirScan::Purgeable { - empty_dirs, + OrphanDirScan::Unreadable => return Ok(false), + OrphanDirScan::Scanned { + blocked: disk_blocked, + dirs, committed_files, } => { - existed = true; - per_disk_dirs.push((i, empty_dirs, committed_files)); + blocked.extend(disk_blocked); + per_disk.push((i, dirs, committed_files)); } OrphanDirScan::Missing => {} } } - - if !existed { - return Ok(false); + if !blocked.is_empty() { + // Whatever is purgeable goes now; what blocks the rest will still + // block it on the next empty listing, so do not rescan for a while. + self.record_orphan_purge_backoff(backoff_key); } - // Phase 2: remove only the files classified as committed residue, then - // remove directories children-first. Every directory delete is - // non-recursive, so a directory that concurrently gained an object fails - // with DirectoryNotEmpty and is skipped — a racing PutObject is never + // Phase 2: remove only the files classified as committed residue in + // unblocked data dirs, then remove unblocked directories + // children-first. Every directory delete is non-recursive, so a + // directory that concurrently gained an object fails with + // DirectoryNotEmpty and is skipped — a racing PutObject is never // clobbered. - for (i, empty_dirs, committed_files) in per_disk_dirs { + let mut purged = false; + for (i, dirs, committed_files) in per_disk { let Some(disk) = disks[i].as_ref() else { continue }; - Self::delete_purgeable_orphan_entries(disk, bucket, object, empty_dirs, committed_files).await; + let dirs = dirs.into_iter().filter(|dir| !blocked.contains(dir)).collect::>(); + let committed_files = committed_files + .into_iter() + .filter(|file| { + file.rsplit_once(SLASH_SEPARATOR) + .is_some_and(|(data_dir, _)| !blocked.contains(data_dir)) + }) + .collect::>(); + if dirs.is_empty() && committed_files.is_empty() { + continue; + } + purged = true; + Self::delete_purgeable_orphan_entries(disk, bucket, object, dirs, committed_files).await; } - Ok(true) + Ok(purged) + } + + fn orphan_purge_in_backoff(&self, key: &str) -> bool { + let backoff = self + .orphan_purge_backoff + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + backoff + .get(key) + .is_some_and(|scanned_at| scanned_at.elapsed() < ORPHAN_PURGE_BACKOFF) + } + + fn record_orphan_purge_backoff(&self, key: String) { + let mut backoff = self + .orphan_purge_backoff + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + let now = Instant::now(); + backoff.retain(|_, scanned_at| now.duration_since(*scanned_at) < ORPHAN_PURGE_BACKOFF); + if backoff.len() >= ORPHAN_PURGE_BACKOFF_MAX_ENTRIES + && let Some(oldest) = backoff + .iter() + .min_by_key(|(_, scanned_at)| **scanned_at) + .map(|(key, _)| key.clone()) + { + backoff.remove(&oldest); + } + backoff.insert(key, now); + } + + #[cfg(test)] + pub(crate) fn clear_orphan_purge_backoff(&self) { + self.orphan_purge_backoff + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .clear(); } /// Reclaim orphaned physical data directories under `bucket/object` that no @@ -7621,13 +7711,15 @@ mod tests { .await .expect("committed delete marker should be written"); - let OrphanDirScan::Purgeable { - empty_dirs, + let OrphanDirScan::Scanned { + blocked, + dirs: empty_dirs, committed_files, } = SetDisks::scan_orphan_dir(&disk, "bucket", "pfx/").await else { panic!("committed residue should be classified as purgeable"); }; + assert!(blocked.is_empty(), "committed residue alone blocks nothing"); let nested_object = residue.join("nested"); tokio::fs::create_dir_all(&nested_object) @@ -7669,13 +7761,15 @@ mod tests { .await .expect("committed delete marker should be written"); - let OrphanDirScan::Purgeable { - empty_dirs, + let OrphanDirScan::Scanned { + blocked, + dirs: empty_dirs, committed_files, } = SetDisks::scan_orphan_dir(&disk, "bucket", "pfx/").await else { panic!("committed residue should be classified as purgeable"); }; + assert!(blocked.is_empty(), "committed residue alone blocks nothing"); tokio::fs::set_permissions(&residue, std::fs::Permissions::from_mode(0o555)) .await .expect("residue directory should become read-only"); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 0fd0cb6be..f75a0a177 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -3900,6 +3900,12 @@ pub struct SetDisks { /// writes skip the global registry mutex (backlog#1315). `Arc` so clones of /// a set share one generation marker. capacity_dirty_generation: Arc, + /// Orphan prefixes whose last purge scan met data that can never be + /// purged by listing (residue without a committed marker, an in-flight + /// write), keyed by `bucket/prefix` with the time of that scan. Empty + /// listings of such a prefix are frequent and the scan reads the whole + /// subtree, so it is not repeated within `ORPHAN_PURGE_BACKOFF` (#6898). + orphan_purge_backoff: Arc>>, #[cfg(test)] storage_class_config_override: Arc>>>, #[cfg(test)] @@ -4690,6 +4696,7 @@ impl SetDisks { ctx, capacity_scope_cache: Arc::new(std::sync::RwLock::new(CapacityScopeCache::default())), capacity_dirty_generation: Arc::new(AtomicU64::new(u64::MAX)), + orphan_purge_backoff: Arc::new(std::sync::Mutex::new(HashMap::new())), #[cfg(test)] storage_class_config_override: Arc::new(std::sync::RwLock::new(None)), #[cfg(test)] @@ -9337,6 +9344,175 @@ mod tests { ); } + async fn write_committed_residue(object_dir: &std::path::Path) -> std::path::PathBuf { + let residue = object_dir.join(Uuid::new_v4().to_string()); + fs::create_dir_all(&residue) + .await + .expect("committed data directory should be created"); + fs::write(residue.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + fs::write( + residue.join(format!("{}{}", crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX, Uuid::new_v4())), + [], + ) + .await + .expect("committed delete marker should be written"); + residue + } + + // #6898: residue from a build that never wrote delete markers blocks only + // its own ancestor chain; a committed sibling subtree is still reclaimed + // and the call reports the partial purge. + #[tokio::test] + async fn purge_orphan_dir_object_reclaims_committed_subtree_beside_blocked_subtree() { + let (dir0, disk0) = make_single_local_disk().await; + let (dir1, disk1) = make_single_local_disk().await; + + let committed = write_committed_residue(&dir0.path().join("bucket/pfx/a/obj")).await; + let unmarked = dir1.path().join("bucket/pfx/b/obj").join(Uuid::new_v4().to_string()); + fs::create_dir_all(&unmarked) + .await + .expect("unmarked residue should be created"); + fs::write(unmarked.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + fs::create_dir_all(dir0.path().join("bucket/pfx/b/obj")) + .await + .expect("disk0 copy of the blocked chain should be created"); + + let set = make_set_disks_with(vec![Some(disk0), Some(disk1)]).await; + let purged = set + .purge_orphan_dir_object("bucket", "pfx/") + .await + .expect("scan should succeed"); + + assert!(purged, "the committed subtree must be reported as purged"); + assert!(!committed.exists(), "committed residue beside a blocked subtree must be reclaimed"); + assert!(!dir0.path().join("bucket/pfx/a").exists(), "the reclaimed subtree's directories must go"); + assert!(unmarked.join("part.1").exists(), "unmarked residue must survive"); + assert!( + dir0.path().join("bucket/pfx/b/obj").exists(), + "a directory blocked on another disk must be left alone on this one" + ); + assert!(dir0.path().join("bucket/pfx").exists() && dir1.path().join("bucket/pfx").exists()); + } + + // A data dir that is committed residue on one disk but still holds an + // object below it on another disk is blocked on every disk. + #[tokio::test] + async fn purge_orphan_dir_object_blocks_committed_files_by_other_disks_data() { + let (dir0, disk0) = make_single_local_disk().await; + let (dir1, disk1) = make_single_local_disk().await; + + let data_dir = Uuid::new_v4(); + let committed = dir0.path().join("bucket/pfx/x").join(data_dir.to_string()); + fs::create_dir_all(&committed) + .await + .expect("committed data directory should be created"); + fs::write(committed.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + fs::write( + committed.join(format!("{}{}", crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX, Uuid::new_v4())), + [], + ) + .await + .expect("committed delete marker should be written"); + let nested = dir1.path().join("bucket/pfx/x").join(data_dir.to_string()).join("nested"); + fs::create_dir_all(&nested) + .await + .expect("nested object dir should be created"); + fs::write(nested.join(STORAGE_FORMAT_FILE), b"meta") + .await + .expect("nested object metadata should be written"); + + let set = make_set_disks_with(vec![Some(disk0), Some(disk1)]).await; + let purged = set + .purge_orphan_dir_object("bucket", "pfx/") + .await + .expect("scan should succeed"); + + assert!(!purged); + assert!( + committed.join("part.1").exists(), + "committed files under a dir blocked elsewhere must remain" + ); + assert!(nested.join(STORAGE_FORMAT_FILE).exists()); + } + + // An unreadable directory anywhere under the prefix aborts the purge on + // every disk before anything is deleted. + #[cfg(unix)] + #[tokio::test] + async fn purge_orphan_dir_object_aborts_when_a_directory_is_unreadable() { + use std::os::unix::fs::PermissionsExt; + + let (dir, disk) = make_single_local_disk().await; + let committed = write_committed_residue(&dir.path().join("bucket/pfx/b/obj")).await; + let sealed = dir.path().join("bucket/pfx/a"); + fs::create_dir_all(&sealed).await.expect("sealed directory should be created"); + fs::set_permissions(&sealed, std::fs::Permissions::from_mode(0o000)) + .await + .expect("sealed directory should become unreadable"); + + let set = make_set_disks_with(vec![Some(disk)]).await; + let purged = set.purge_orphan_dir_object("bucket", "pfx/").await; + fs::set_permissions(&sealed, std::fs::Permissions::from_mode(0o755)) + .await + .expect("sealed directory should be restored"); + + assert!(matches!(purged, Ok(false)), "an unreadable directory must abort the purge"); + assert!(committed.join("part.1").exists(), "nothing may be deleted once classification failed"); + } + + // After a scan met unpurgeable data the prefix is not rescanned for a + // while, even when new committed residue appears under it. + #[tokio::test] + async fn purge_orphan_dir_object_backs_off_after_blocked_scan() { + let (dir, disk) = make_single_local_disk().await; + let unmarked = dir.path().join("bucket/pfx/b/obj").join(Uuid::new_v4().to_string()); + fs::create_dir_all(&unmarked) + .await + .expect("unmarked residue should be created"); + fs::write(unmarked.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + + let set = make_set_disks_with(vec![Some(disk)]).await; + assert!( + !set.purge_orphan_dir_object("bucket", "pfx/") + .await + .expect("scan should succeed") + ); + + let committed = write_committed_residue(&dir.path().join("bucket/pfx/a/obj")).await; + assert!( + !set.purge_orphan_dir_object("bucket", "pfx/") + .await + .expect("scan should succeed"), + "a prefix in backoff is not rescanned" + ); + assert!(committed.exists()); + assert!( + set.purge_orphan_dir_object("bucket", "pfx/a/") + .await + .expect("scan should succeed"), + "backoff is per prefix, a narrower prefix still purges" + ); + assert!(!committed.exists()); + + let committed = write_committed_residue(&dir.path().join("bucket/pfx/a/obj")).await; + set.clear_orphan_purge_backoff(); + assert!( + set.purge_orphan_dir_object("bucket", "pfx/") + .await + .expect("scan should succeed") + ); + assert!(!committed.exists()); + assert!(unmarked.join("part.1").exists()); + } + // Build an `xl.meta` under `object_dir` whose versions reference `data_dirs` // (one Object version per data dir, each with its own version id). async fn write_object_meta_with_data_dirs(object_dir: &std::path::Path, bucket: &str, object: &str, data_dirs: &[Uuid]) { diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 7f2508e27..9a1de1cb8 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -9678,6 +9678,131 @@ mod test { } } + #[tokio::test] + async fn empty_delimiter_listing_of_residue_ancestor_hides_and_purges_whole_tree() { + 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 = "listing-purge-ancestor-bucket"; + init_bucket_metadata_sys(store.clone(), Vec::new()).await; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created with authoritative metadata"); + let data_dir = uuid::Uuid::new_v4(); + let transaction = uuid::Uuid::new_v4(); + for dir in &dirs { + let residue = dir + .path() + .join(bucket) + .join("metrics/kubelet/2026/08/28/23/74992556388248657933757.parquet") + .join(data_dir.to_string()); + tokio::fs::create_dir_all(&residue) + .await + .expect("committed delete residue should be created"); + tokio::fs::write(residue.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + tokio::fs::write( + residue.join(format!("{}{}", crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX, transaction)), + [], + ) + .await + .expect("committed delete marker should be written"); + } + + // Browsing an ancestor of the deleted key must not show the empty + // date folders, and the empty result reclaims the whole residue tree + // in one pass instead of one level per listing. + let result = store + .clone() + .list_objects_generic(bucket, "metrics/", None, Some("/".to_owned()), 1000, false) + .await + .expect("delimiter listing should succeed"); + assert!(result.objects.is_empty()); + assert!(result.prefixes.is_empty(), "ancestors of delete residue must not surface as prefixes"); + for dir in &dirs { + assert!( + !dir.path().join(bucket).join("metrics").exists(), + "the empty delimiter listing should reclaim the committed delete residue tree under it" + ); + } + } + + #[tokio::test] + async fn empty_delimiter_listing_of_mixed_residue_ancestor_reclaims_only_committed_subtree() { + 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 = "listing-purge-mixed-bucket"; + init_bucket_metadata_sys(store.clone(), Vec::new()).await; + store + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created with authoritative metadata"); + let committed_dir = uuid::Uuid::new_v4(); + let unmarked_dir = uuid::Uuid::new_v4(); + let transaction = uuid::Uuid::new_v4(); + for dir in &dirs { + let committed = dir + .path() + .join(bucket) + .join("metrics/cpu/2026/08/28/22/a.parquet") + .join(committed_dir.to_string()); + tokio::fs::create_dir_all(&committed) + .await + .expect("committed delete residue should be created"); + tokio::fs::write(committed.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + tokio::fs::write( + committed.join(format!("{}{}", crate::disk::local::DELETE_DATA_DIR_MARKER_PREFIX, transaction)), + [], + ) + .await + .expect("committed delete marker should be written"); + + // Residue left by a build that never wrote delete markers, or a + // PUT still streaming its parts: indistinguishable, so never purged. + let unmarked = dir + .path() + .join(bucket) + .join("metrics/kubelet/2026/08/28/23/74992556388248657933757.parquet") + .join(unmarked_dir.to_string()); + tokio::fs::create_dir_all(&unmarked) + .await + .expect("unmarked residue should be created"); + tokio::fs::write(unmarked.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + } + + let result = store + .clone() + .list_objects_generic(bucket, "metrics/", None, Some("/".to_owned()), 1000, false) + .await + .expect("delimiter listing should succeed"); + assert!(result.objects.is_empty()); + assert!(result.prefixes.is_empty(), "neither residue subtree may surface as a prefix"); + for dir in &dirs { + let metrics = dir.path().join(bucket).join("metrics"); + assert!( + !metrics.join("cpu").exists(), + "the committed subtree must be reclaimed even though a sibling subtree is blocked" + ); + assert!( + metrics + .join("kubelet/2026/08/28/23/74992556388248657933757.parquet") + .join(unmarked_dir.to_string()) + .join("part.1") + .exists(), + "residue without a committed marker must survive the purge" + ); + } + } + #[test] fn list_objects_index_provider_state_uses_lifecycle_active_generation() { let provider = ListObjectsIndexProviderState::walker_key_only();