diff --git a/crates/ecstore/src/set_disk/core/io_primitives.rs b/crates/ecstore/src/set_disk/core/io_primitives.rs index f9de75a9a..6f2d109c6 100644 --- a/crates/ecstore/src/set_disk/core/io_primitives.rs +++ b/crates/ecstore/src/set_disk/core/io_primitives.rs @@ -2713,51 +2713,158 @@ impl SetDisks { versions } - #[allow(dead_code)] + /// Reclaim the old (now dereferenced) `object/` on the disks + /// that just committed the new version. + /// + /// # Deliberate divergence from MinIO (backlog#898) + /// + /// This runs *after* the write is authoritatively committed (`rename_data` + /// returned `Ok`, i.e. the new version is durable on >= write_quorum disks + /// and immediately readable). The target `object/` has already + /// been dereferenced from every committed replica's `xl.meta`, so removing + /// it is pure space reclamation. MinIO couples a below-quorum failure here + /// back into the client response (`erasure-object.go:1577`); RustFS + /// historically mirrored that (`reduce_write_quorum_errs` -> `Err` -> + /// 503/SlowDown), producing a **false-negative ACK** for an already-durable + /// write. We deliberately break that coupling: this function **never returns + /// `Err`**. It returns a structured [`OldDataDirCleanup`] receipt whose + /// fields are signals only — none of them can negate an already-ACKed write. + /// + /// The disk-health signal that MinIO raises via 503 is not dropped; the + /// caller re-surfaces it by enqueuing an object heal on residue (see + /// `report_old_data_dir_cleanup`), and the leaked residue is made observable + /// via `rustfs_old_data_dir_leaked_total`. #[tracing::instrument(level = "debug", skip(self, disks))] pub(in crate::set_disk) async fn commit_rename_data_dir( &self, disks: &[Option], bucket: &str, object: &str, - data_dir: &str, + old_data_dir: &str, + committed_data_dir: &str, write_quorum: usize, - ) -> disk::error::Result<()> { - let file_path = Arc::new(format!("{object}/{data_dir}")); + ) -> OldDataDirCleanup { + // Anti-misdelete guard (parity retained): MinIO sets `res.OldDataDir=""` + // when old == new (`xl-storage.go:2796`) and `commitRenameDataDir` skips + // the empty value (`:1837`). Rather miss a reclaim than delete the data + // dir we just committed. `#864` isolation: never touch the commit point. + if old_data_dir.is_empty() || old_data_dir == committed_data_dir { + return OldDataDirCleanup::default(); + } + + let file_path = Arc::new(format!("{object}/{old_data_dir}")); let bucket = Arc::new(bucket.to_string()); - let futures = disks.iter().map(|disk| { + // A disk slot was actually targeted for deletion iff it is `Some`; the + // `None` slots are ignored placeholders and must be excluded from the + // attempted/reclaimed/residue accounting (they leak nothing). + let attempted: Vec = disks.iter().map(|d| d.is_some()).collect(); + // Test-only object filter for the delete fault-injection seam. In + // non-test builds `cleanup_injected_error` is a `None`-returning no-op, + // so this clone and the per-disk check compile away to nothing. + let object_for_fault = Arc::new(object.to_string()); + + let futures = disks.iter().enumerate().map(|(idx, disk)| { let file_path = file_path.clone(); let bucket = bucket.clone(); let disk = disk.clone(); + let object_for_fault = object_for_fault.clone(); tokio::spawn(async move { + if let Some(err) = Self::cleanup_injected_error(&object_for_fault, idx) { + return Some(err); + } if let Some(disk) = disk { - (disk - .delete( - &bucket, - &file_path, - DeleteOptions { - recursive: true, - ..Default::default() - }, - ) - .await) - .err() + disk.delete( + &bucket, + &file_path, + DeleteOptions { + recursive: true, + ..Default::default() + }, + ) + .await + .err() } else { + // `None` slot: ignored placeholder. It is not `attempted`, so + // classification excludes it from residue regardless. Some(DiskError::DiskNotFound) } }) }); - let errs: Vec> = join_all(futures) - .await - .into_iter() - .map(|e| e.unwrap_or(Some(DiskError::Unexpected))) - .collect(); + let errs: Vec> = join_all(futures).await.into_iter().map(map_cleanup_join_result).collect(); - if let Some(err) = reduce_write_quorum_errs(&errs, OBJECT_OP_IGNORED_ERRS, write_quorum) { - return Err(err); + classify_old_data_dir_cleanup(&errs, &attempted, write_quorum) + } + + /// Test-only fault-injection seam for the old-data-dir cleanup path + /// (backlog#898 §5). In production this is inlined to `None` and adds no + /// behavior; only the `#[cfg(test)]` variant consults the fault registry. + #[cfg(test)] + fn cleanup_injected_error(object: &str, disk_index: usize) -> Option { + cleanup_fault_injection::injected_error(object, disk_index) + } + + #[cfg(not(test))] + #[inline(always)] + fn cleanup_injected_error(_object: &str, _disk_index: usize) -> Option { + None + } + + /// Report a post-commit old-data-dir cleanup receipt: emit metrics, warn on + /// residue/below-quorum, and — on residue — enqueue an object heal. + /// + /// Open-question defaults adopted from backlog#898 §7 (flagged for + /// maintainer confirmation): + /// - the heal is enqueued **inline** (awaited), matching the existing + /// `multipart.rs` fire-and-forget precedent; it runs only on the rare + /// residue path, after the object write lock is released, and is + /// deduplicated/back-pressured by heal admission — so the ACK tail impact + /// is negligible. (Alternative: `tokio::spawn`.) + /// - the heal is enqueued strictly on `has_residue()`, not on + /// `below_quorum` alone (a pure parity lens that can flag offline/`None` + /// slots which leak nothing). + pub(in crate::set_disk) async fn report_old_data_dir_cleanup( + &self, + bucket: &str, + object: &str, + old_dir: &str, + c: &OldDataDirCleanup, + ) { + let actions = old_data_dir_cleanup_actions(c); + + rustfs_io_metrics::record_old_data_dir_cleanup(c.attempted, c.reclaimed, c.unreclaimed_disks.len(), c.below_quorum); + + if actions.warn { + warn!( + component = LOG_COMPONENT_ECSTORE, + subsystem = LOG_SUBSYSTEM_SET_DISK, + bucket = %bucket, + object = %object, + old_dir = %old_dir, + attempted = c.attempted, + reclaimed = c.reclaimed, + unreclaimed = ?c.unreclaimed_disks, + below_quorum = c.below_quorum, + "old data dir cleanup left residue after committed write" + ); } - Ok(()) + if actions.enqueue_heal { + // Disk-health signal (replacing MinIO's 503) + convergence hook: the + // leaked local data dir is physically reclaimed by heal_object -> + // reclaim_orphan_data_dirs. Reuses the existing heal channel, which + // deduplicates and back-pressures via admission; failures only drop + // the return value (same shape as multipart's existing heal enqueue). + let _ = + rustfs_common::heal_channel::send_heal_request(rustfs_common::heal_channel::create_heal_request_with_options( + bucket.to_string(), + Some(object.to_string()), + false, + Some(rustfs_common::heal_channel::HealChannelPriority::Normal), + Some(self.pool_index), + Some(self.set_index), + )) + .await; + } } #[tracing::instrument(skip(self))] @@ -3545,6 +3652,160 @@ impl SetDisks { } } +/// Structured receipt for post-commit old-data-dir reclamation (backlog#898). +/// +/// Every field is a *signal*; none of them can negate an already-ACKed write. +/// `commit_rename_data_dir` returns this by value (never a `Result`), so the two +/// call sites cannot `?`-propagate a 503 out of a successful commit. +#[derive(Debug, Clone, PartialEq, Eq, Default)] +pub(in crate::set_disk) struct OldDataDirCleanup { + /// Number of disks a delete was actually issued to (`disks[i].is_some()`). + pub attempted: usize, + /// Number of attempted disks that returned `Ok` or a not-found variant + /// (a missing dir == already reclaimed). + pub reclaimed: usize, + /// Indices of attempted disks that failed with a non-ignored, non-not-found + /// error (including task panic/cancel). This is the residue that actually + /// leaks and drives the leak metric + heal enqueue. + pub unreclaimed_disks: Vec, + /// Parity-lens signal only: `reduce_write_quorum_errs` over the raw per-disk + /// errors. It can be `true` while `unreclaimed_disks` is empty (e.g. many + /// offline/`None` slots, which are ignored). Used for logging/metrics only — + /// it never participates in any return decision. + pub below_quorum: bool, +} + +impl OldDataDirCleanup { + /// Whether any attempted disk left a real, leaked residue. + pub fn has_residue(&self) -> bool { + !self.unreclaimed_disks.is_empty() + } +} + +/// not-found normalized to success (parity with MinIO `commitRenameDataDir`). +fn is_cleanup_not_found(e: &DiskError) -> bool { + matches!(e, DiskError::FileNotFound | DiskError::VolumeNotFound | DiskError::PathNotFound) +} + +/// Map a tokio join result into a reducible per-disk error. +/// +/// A task panic/cancel is mapped to a **non-ignored** `DiskError::other`, never +/// normalized to `DiskNotFound`: a panic is not a "disk absent" condition and +/// must not be silently swallowed as an ignorable error (fixes the historical +/// `Unexpected`/`DiskNotFound` misclassification). +fn map_cleanup_join_result(joined: std::result::Result, tokio::task::JoinError>) -> Option { + match joined { + Ok(res) => res, + Err(join_err) => Some(DiskError::other(format!("old data dir cleanup task failed: {join_err}"))), + } +} + +/// Pure classification: reduce per-disk cleanup results into a receipt. +/// +/// `errs[i] == None` means success; `attempted[i]` marks that `disks[i]` was a +/// real target (`Some`). Untargeted `None` slots still contribute to the +/// `below_quorum` parity lens (via their ignored placeholder error) but are +/// excluded from attempted/reclaimed/residue. +fn classify_old_data_dir_cleanup(errs: &[Option], attempted: &[bool], write_quorum: usize) -> OldDataDirCleanup { + debug_assert_eq!(errs.len(), attempted.len()); + let below_quorum = reduce_write_quorum_errs(errs, OBJECT_OP_IGNORED_ERRS, write_quorum).is_some(); + let mut reclaimed = 0usize; + let mut attempted_count = 0usize; + let mut unreclaimed_disks = Vec::new(); + for (i, err) in errs.iter().enumerate() { + if !attempted.get(i).copied().unwrap_or(false) { + continue; + } + attempted_count += 1; + match err { + None => reclaimed += 1, + Some(e) if is_cleanup_not_found(e) => reclaimed += 1, + Some(_) => unreclaimed_disks.push(i), + } + } + OldDataDirCleanup { + attempted: attempted_count, + reclaimed, + unreclaimed_disks, + below_quorum, + } +} + +/// Decision derived purely from a cleanup receipt, kept separate from `&self` +/// so the policy is unit-testable (backlog#898 §6). +#[derive(Debug, Default, PartialEq, Eq)] +pub(in crate::set_disk) struct CleanupActions { + pub warn: bool, + pub emit_leak_metric: bool, + pub enqueue_heal: bool, +} + +/// Decide what to do about a cleanup receipt. Residue is the only trigger for +/// the leak metric and the heal enqueue; `below_quorum` only widens the warn. +fn old_data_dir_cleanup_actions(c: &OldDataDirCleanup) -> CleanupActions { + CleanupActions { + warn: c.has_residue() || c.below_quorum, + emit_leak_metric: c.has_residue(), + enqueue_heal: c.has_residue(), + } +} + +/// Test-only delete fault-injection seam for the old-data-dir cleanup path +/// (backlog#898 §5). +/// +/// This entire module is `#[cfg(test)]`, so it never compiles into production. +/// Because the cleanup deletes run in spawned tasks (potentially on different +/// worker threads), the registry is a process-global keyed by object name; a +/// test scopes its injection to a unique object name via [`fail_cleanup_on`], +/// and the returned guard clears the registry on drop. +#[cfg(test)] +pub(in crate::set_disk) mod cleanup_fault_injection { + use super::DiskError; + use std::sync::{Mutex, OnceLock}; + + #[derive(Default)] + struct FaultState { + object: Option, + fail_indices: Vec, + } + + fn state() -> &'static Mutex { + static STATE: OnceLock> = OnceLock::new(); + STATE.get_or_init(|| Mutex::new(FaultState::default())) + } + + /// Guard that clears the fault registry on drop. + pub struct FaultGuard; + + impl Drop for FaultGuard { + fn drop(&mut self) { + if let Ok(mut s) = state().lock() { + *s = FaultState::default(); + } + } + } + + /// Force old-data-dir cleanup deletes for `object` on the given disk indices + /// to fail with a transient error, until the returned guard is dropped. + #[must_use] + pub fn fail_cleanup_on(object: &str, indices: &[usize]) -> FaultGuard { + let mut s = state().lock().expect("cleanup fault registry poisoned"); + s.object = Some(object.to_string()); + s.fail_indices = indices.to_vec(); + FaultGuard + } + + pub(super) fn injected_error(object: &str, disk_index: usize) -> Option { + let s = state().lock().expect("cleanup fault registry poisoned"); + match &s.object { + Some(target) if target == object && s.fail_indices.contains(&disk_index) => { + Some(DiskError::other("injected old-data-dir cleanup failure (test-only)")) + } + _ => None, + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -3823,6 +4084,122 @@ mod tests { release_read_repair_heal_reservation(&key).await; } + // ======================================================================== + // backlog#898 — groups A/C: pure old-data-dir cleanup classification. + // ======================================================================== + use crate::disk::error_reduce::is_ignored_err; + + // A1: all deletes succeed => full reclaim, no residue, not below quorum. + #[test] + fn cleanup_receipt_all_deletes_succeed_marks_full_reclaim() { + let errs = vec![None, None, None]; + let attempted = vec![true, true, true]; + let r = classify_old_data_dir_cleanup(&errs, &attempted, 2); + assert_eq!(r.attempted, 3); + assert_eq!(r.reclaimed, 3); + assert!(r.unreclaimed_disks.is_empty()); + assert!(!r.below_quorum); + assert!(!r.has_residue()); + } + + // A2: not-found normalized to success (parity with MinIO commitRenameDataDir). + #[test] + fn cleanup_receipt_not_found_counts_as_reclaimed() { + let errs = vec![None, Some(DiskError::FileNotFound), Some(DiskError::VolumeNotFound)]; + let attempted = vec![true, true, true]; + let r = classify_old_data_dir_cleanup(&errs, &attempted, 1); + assert_eq!(r.reclaimed, 3, "absent dir must count as already reclaimed"); + assert!(r.unreclaimed_disks.is_empty()); + } + + // A3: non-ignored transient failures below quorum => a receipt (not an Err), + // below_quorum true, residue set listing exactly the failed disks. + #[test] + fn cleanup_receipt_transient_failures_are_residue_not_error() { + let errs = vec![ + None, + None, + Some(DiskError::other("connection reset")), + Some(DiskError::other("io timeout")), + ]; + let attempted = vec![true, true, true, true]; + let r = classify_old_data_dir_cleanup(&errs, &attempted, 3); + assert_eq!(r.attempted, 4); + assert_eq!(r.reclaimed, 2); + assert_eq!(r.unreclaimed_disks, vec![2, 3], "residue set must be exact"); + assert!(r.below_quorum, "2 achieved < wq 3"); + assert!(r.has_residue()); + } + + // A4: a join error (panic/cancel) must not be masked as DiskNotFound; it must + // land in the residue and never be treated as an ignored error. Directly + // regresses the old `Unexpected`/`DiskNotFound` misclassification. + #[tokio::test] + async fn cleanup_join_error_maps_to_non_ignored_error_not_disk_not_found() { + let handle: tokio::task::JoinHandle> = tokio::spawn(async { panic!("boom") }); + let joined = handle.await; + assert!(joined.is_err(), "task must have panicked"); + let mapped = map_cleanup_join_result(joined).expect("panic must surface as an error"); + assert_ne!( + mapped, + DiskError::DiskNotFound, + "a task panic must never be masked as an ignorable DiskNotFound" + ); + assert!( + !is_ignored_err(OBJECT_OP_IGNORED_ERRS, &mapped), + "a task panic must not be treated as an ignored error" + ); + let r = classify_old_data_dir_cleanup(&[Some(mapped)], &[true], 1); + assert_eq!(r.unreclaimed_disks, vec![0]); + assert!(r.has_residue()); + } + + // A6: all-None targeted slots => no attempt, no residue; below_quorum may be + // true (parity lens) but nothing actually leaks. Pins the two as independent. + #[test] + fn cleanup_receipt_ignores_untargeted_none_slots() { + let errs = vec![Some(DiskError::DiskNotFound), Some(DiskError::DiskNotFound)]; + let attempted = vec![false, false]; + let r = classify_old_data_dir_cleanup(&errs, &attempted, 1); + assert_eq!(r.attempted, 0); + assert_eq!(r.reclaimed, 0); + assert!(r.unreclaimed_disks.is_empty()); + assert!(!r.has_residue(), "None-only slots never leak anything"); + assert!(r.below_quorum, "parity lens still flags it, but nothing is actually leaked"); + } + + // divergence: residue and below_quorum are independent signals — quorum met + // can still leak. + #[test] + fn cleanup_receipt_residue_and_below_quorum_are_independent_signals() { + let errs = vec![None, None, None, Some(DiskError::other("io"))]; + let attempted = vec![true, true, true, true]; + let r = classify_old_data_dir_cleanup(&errs, &attempted, 3); + assert!(!r.below_quorum, "3 achieved >= wq 3"); + assert_eq!(r.unreclaimed_disks, vec![3]); + assert!(r.has_residue(), "leak metric must fire even when write quorum was met"); + } + + // C: persistent failure => leak metric + heal enqueue decided, no Err/panic. + #[test] + fn persistent_cleanup_failure_triggers_leak_metric_and_heal_not_error() { + let errs = vec![None, Some(DiskError::other("disk down")), Some(DiskError::other("disk down"))]; + let attempted = vec![true, true, true]; + let receipt = classify_old_data_dir_cleanup(&errs, &attempted, 2); + assert!(receipt.has_residue()); + let actions = old_data_dir_cleanup_actions(&receipt); + assert!(actions.emit_leak_metric, "persistent residue must be recorded by an observable metric"); + assert!(actions.enqueue_heal, "persistent residue must enqueue a heal (disk-health signal)"); + assert!(actions.warn); + } + + // A clean receipt takes no action. + #[test] + fn clean_cleanup_receipt_triggers_no_actions() { + let receipt = classify_old_data_dir_cleanup(&[None, None], &[true, true], 1); + assert_eq!(old_data_dir_cleanup_actions(&receipt), CleanupActions::default()); + } + // ------------------------------------------------------------------ // backlog#900: pick_latest_quorum_files_info must survive a single // corrupt-part disk (even in the merged representative slot) by deriving diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 779f8826e..a00af9e00 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -4667,6 +4667,128 @@ mod tests { assert!(obj0.join(other.to_string()).exists(), "locally referenced data dir must survive"); } + // backlog#898 A5: old == committed data dir => anti-misdelete guard skips the + // whole cleanup; the live (just-committed) dir must be left untouched. + #[tokio::test] + async fn commit_rename_data_dir_skips_delete_when_old_equals_committed_dir() { + let (dir, disk) = make_single_local_disk().await; + let root = dir.path(); + let bucket = "bucket"; + let object = "object"; + fs::create_dir_all(root.join(bucket)) + .await + .expect("bucket volume should be created"); + + let same_dir = Uuid::parse_str("55555555-5555-5555-5555-555555555555").expect("dir should parse"); + let data_path = root.join(bucket).join(object).join(same_dir.to_string()); + fs::create_dir_all(&data_path).await.expect("data dir should be created"); + fs::write(data_path.join("part.1"), b"live") + .await + .expect("live part should be written"); + + let set = make_set_disks_with(vec![Some(disk.clone())]).await; + let cleanup = set + .commit_rename_data_dir(&[Some(disk.clone())], bucket, object, &same_dir.to_string(), &same_dir.to_string(), 1) + .await; + + assert_eq!(cleanup.attempted, 0, "guard must skip: no delete may be issued"); + assert!(cleanup.unreclaimed_disks.is_empty()); + assert!(!cleanup.has_residue()); + assert!(data_path.join("part.1").exists(), "committed data dir must NOT be deleted"); + } + + // backlog#898 A5b: old != committed and old dir exists => normal reclaim; the + // dereferenced old dir is physically removed and the receipt reports success. + #[tokio::test] + async fn commit_rename_data_dir_reclaims_distinct_old_dir() { + let (dir, disk) = make_single_local_disk().await; + let root = dir.path(); + let bucket = "bucket"; + let object = "object"; + fs::create_dir_all(root.join(bucket)) + .await + .expect("bucket volume should be created"); + + let old_dir = Uuid::parse_str("11111111-1111-1111-1111-111111111111").expect("old dir should parse"); + let new_dir = Uuid::parse_str("22222222-2222-2222-2222-222222222222").expect("new dir should parse"); + let old_path = root.join(bucket).join(object).join(old_dir.to_string()); + fs::create_dir_all(&old_path).await.expect("old data dir should be created"); + fs::write(old_path.join("part.1"), b"stale") + .await + .expect("stale part should be written"); + + let set = make_set_disks_with(vec![Some(disk.clone())]).await; + let cleanup = set + .commit_rename_data_dir(&[Some(disk.clone())], bucket, object, &old_dir.to_string(), &new_dir.to_string(), 1) + .await; + + assert_eq!(cleanup.attempted, 1); + assert_eq!(cleanup.reclaimed, 1); + assert!(!cleanup.has_residue()); + assert!(!cleanup.below_quorum); + assert!(!old_path.exists(), "dereferenced old data dir must be physically removed"); + } + + // backlog#898 group B (end-to-end): a real overwrite whose post-commit old + // data dir cleanup is forced to fail (via the test-only fault seam) must + // still return 200 with the new ObjectInfo — no false-negative ACK. + #[tokio::test] + async fn put_object_overwrite_returns_ok_when_old_data_dir_cleanup_fails() { + use crate::set_disk::core::io_primitives::cleanup_fault_injection; + + let set_disks = make_local_bucket_test_set_disks().await; + let bucket = "bucket-cleanup-fault"; + let object = "object-below-quorum-cleanup"; + + set_disks + .make_bucket(bucket, &MakeBucketOptions::default()) + .await + .expect("bucket should be created"); + + // v1: establishes an old data dir that the overwrite will try to reclaim. + let mut reader = PutObjReader::from_vec(b"hello".to_vec()); + set_disks + .put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + no_lock: true, + ..ObjectOptions::default() + }, + ) + .await + .expect("first write should succeed"); + + // Force the old-data-dir cleanup delete to fail on disk 0 during the + // overwrite. rename_data still commits; only the GC of the dereferenced + // old dir fails, which must NOT turn the PUT into a 503. + let _fault = cleanup_fault_injection::fail_cleanup_on(object, &[0]); + + let mut reader = PutObjReader::from_vec(b"goodbye!!".to_vec()); + let oi = set_disks + .put_object( + bucket, + object, + &mut reader, + &ObjectOptions { + no_lock: true, + ..ObjectOptions::default() + }, + ) + .await + .expect("overwrite must return Ok even when old-data-dir cleanup fails"); + + assert_eq!(oi.size, 9, "returned ObjectInfo must reflect the new committed write"); + + // The committed new version must be readable with the new size. + let read_back = set_disks + .get_object_info(bucket, object, &ObjectOptions::default()) + .await + .expect("committed object must be readable after a failed cleanup"); + assert_eq!(read_back.size, 9, "HEAD must observe the new version, not stale metadata"); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn test_acquire_dist_delete_object_locks_batch_succeeds_with_two_healthy_lockers() { diff --git a/crates/ecstore/src/set_disk/ops/multipart.rs b/crates/ecstore/src/set_disk/ops/multipart.rs index 7ae6d3a6f..7ee3a0012 100644 --- a/crates/ecstore/src/set_disk/ops/multipart.rs +++ b/crates/ecstore/src/set_disk/ops/multipart.rs @@ -1357,8 +1357,15 @@ impl crate::storage_api_contracts::multipart::MultipartOperations for SetDisks { .await?; if let Some(old_dir) = op_old_dir { - self.commit_rename_data_dir(&cleanup_disks, bucket, object, &old_dir.to_string(), write_quorum) - .await?; + let committed_dir = fi.data_dir.unwrap_or_default().to_string(); + // backlog#898: best-effort reclaim of the dereferenced old data dir. + // Returns a receipt (never `Err`); a failed GC must not turn an + // already-committed multipart completion into a 503. + let cleanup = self + .commit_rename_data_dir(&cleanup_disks, bucket, object, &old_dir.to_string(), &committed_dir, write_quorum) + .await; + self.report_old_data_dir_cleanup(bucket, object, &old_dir.to_string(), &cleanup) + .await; } if let Some(stage_start) = complete_tail_stage_start { diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 2f2e69db3..244db4c10 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -895,12 +895,21 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks { let mut cleanup_stage_ms: Option = None; if let Some(old_dir) = op_old_dir { + let committed_dir = fi.data_dir.unwrap_or_default().to_string(); let cleanup_stage_start = Instant::now(); - self.commit_rename_data_dir(&cleanup_disks, bucket, object, &old_dir.to_string(), write_quorum) - .await?; + // backlog#898: reclaiming the dereferenced old data dir is + // best-effort and returns a receipt (never `Err`). A failed GC + // here must not negate an already-committed, durable write, so we + // deliberately do NOT `?`-propagate it into a 503. On residue the + // report path emits the leak metric and enqueues a heal. + let cleanup = self + .commit_rename_data_dir(&cleanup_disks, bucket, object, &old_dir.to_string(), &committed_dir, write_quorum) + .await; let cleanup_ms = cleanup_stage_start.elapsed().as_millis() as u64; cleanup_stage_ms = Some(cleanup_ms); rustfs_io_metrics::record_put_object_stage_duration("set_disk_old_data_cleanup", cleanup_ms as f64); + self.report_old_data_dir_cleanup(bucket, object, &old_dir.to_string(), &cleanup) + .await; if (cleanup_ms as u128) >= SET_DISK_COMMIT_TAIL_WARN_THRESHOLD_MS { warn!( event = EVENT_SET_DISK_COMMIT_TAIL_SLOW, diff --git a/crates/io-metrics/src/capacity_metrics.rs b/crates/io-metrics/src/capacity_metrics.rs index f76ea85e7..b7118a57b 100644 --- a/crates/io-metrics/src/capacity_metrics.rs +++ b/crates/io-metrics/src/capacity_metrics.rs @@ -172,3 +172,22 @@ pub fn record_capacity_scan_disk( counter!("rustfs_capacity_scan_disk_partial_errors_total", "disk" => disk.to_owned()).increment(1); } } + +/// Record the outcome of a post-commit old-data-dir cleanup (backlog#898). +/// +/// This is fired once per committed overwrite that had a previous data dir to +/// reclaim. `leaked` is the number of disks whose old data dir could not be +/// removed (a non-ignored, non-not-found failure) and is therefore a residue +/// that leaks disk space until a heal reclaims it. The leak counter is the +/// operator-visible backstop for that residue — see backlog#898 §3.4/§5. +#[inline(always)] +pub fn record_old_data_dir_cleanup(attempted: usize, reclaimed: usize, leaked: usize, below_quorum: bool) { + counter!("rustfs_old_data_dir_cleanup_attempted_total").increment(attempted as u64); + counter!("rustfs_old_data_dir_cleanup_reclaimed_total").increment(reclaimed as u64); + if leaked > 0 { + counter!("rustfs_old_data_dir_leaked_total").increment(leaked as u64); + } + if below_quorum { + counter!("rustfs_old_data_dir_cleanup_below_quorum_total").increment(1); + } +} diff --git a/crates/io-metrics/src/lib.rs b/crates/io-metrics/src/lib.rs index fa84c3d7b..a568e2b61 100644 --- a/crates/io-metrics/src/lib.rs +++ b/crates/io-metrics/src/lib.rs @@ -122,7 +122,7 @@ pub use capacity_metrics::{ record_capacity_refresh_joiner, record_capacity_refresh_request, record_capacity_refresh_result, record_capacity_refresh_scope, record_capacity_scan_disk, record_capacity_scan_mode, record_capacity_scan_sampling, record_capacity_stall_detected, record_capacity_symlink, record_capacity_timeout_fallback, record_capacity_update_completed, - record_capacity_update_failed, record_capacity_write_operation, + record_capacity_update_failed, record_capacity_write_operation, record_old_data_dir_cleanup, }; // I/O metrics exports