From 72d76dc9ca3344a4959991958c6e019d3f85cf66 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Thu, 9 Jul 2026 02:22:09 +0800 Subject: [PATCH] fix(object-capacity): make dirty-disk lifecycle race-free and prune ghost entries (#4550) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two dirty-mark lifecycle gaps (S19+S30): - A commit cleared every dirty mark for the disks it scanned, even marks recorded while the walker was already past the written prefix, so the next dirty-subset refresh served stale cached bytes as exact until the next full scan. Dirty marks now carry the instant they were last recorded and a commit only clears marks that predate its scan start (CapacityUpdate::scan_started_at, set by refresh_capacity_with_scope). - The write side marks every disk of an EC set dirty, including remote peers, while scans and clearing only cover local disks — remote or removed disks stayed marked forever, keeping the dirty-disk gauge permanently non-zero with bounded memory stuck behind it. The refresh disk selection now prunes dirty entries outside the current local topology before consulting them. Ref: rustfs/backlog#1020 (S19+S30 from audit rustfs/backlog#1010) --- .../object-capacity/src/capacity_manager.rs | 137 +++++++++++++++++- crates/object-capacity/src/scan.rs | 12 +- 2 files changed, 141 insertions(+), 8 deletions(-) diff --git a/crates/object-capacity/src/capacity_manager.rs b/crates/object-capacity/src/capacity_manager.rs index 3e427bcc3..e960cc361 100644 --- a/crates/object-capacity/src/capacity_manager.rs +++ b/crates/object-capacity/src/capacity_manager.rs @@ -352,6 +352,10 @@ pub struct CapacityUpdate { pub replaces_disk_cache: bool, /// Dirty disks that can be cleared after the update is committed. pub clear_dirty_disks: Vec, + /// When the scan behind this update started; dirty marks recorded after + /// this instant survive the commit (backlog#1020 S19). `None` clears + /// unconditionally. + pub scan_started_at: Option, } impl CapacityUpdate { @@ -366,6 +370,7 @@ impl CapacityUpdate { expected_disk_count: None, replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, } } @@ -380,6 +385,7 @@ impl CapacityUpdate { expected_disk_count: None, replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, } } @@ -394,6 +400,7 @@ impl CapacityUpdate { expected_disk_count: None, replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, } } } @@ -631,8 +638,10 @@ pub struct HybridCapacityManager { cache: Arc>>, /// Write record write_record: Arc>, - /// Dirty disks recorded from write-side scope propagation. - dirty_disks: Arc>>, + /// Dirty disks recorded from write-side scope propagation, keyed to the + /// instant they were last marked so a commit only clears marks that + /// predate its scan (backlog#1020 S19). + dirty_disks: Arc>>, /// Per-disk cache populated after a successful full refresh and updated by dirty subset refreshes. disk_cache: Arc>>, /// Whether the per-disk cache currently covers all known disks. @@ -650,8 +659,9 @@ impl HybridCapacityManager { return; } + let now = Instant::now(); let mut dirty_disks = self.dirty_disks.write().await; - dirty_disks.extend(scopes); + dirty_disks.extend(scopes.into_iter().map(|disk| (disk, now))); record_capacity_dirty_disk_count(dirty_disks.len()); } @@ -670,7 +680,7 @@ impl HybridCapacityManager { write_count: 0, write_buckets: [WriteBucket::default(); WRITE_WINDOW_BUCKETS], })), - dirty_disks: Arc::new(RwLock::new(HashSet::new())), + dirty_disks: Arc::new(RwLock::new(HashMap::new())), disk_cache: Arc::new(RwLock::new(HashMap::new())), disk_cache_complete: Arc::new(RwLock::new(false)), config, @@ -765,7 +775,18 @@ impl HybridCapacityManager { if !update.clear_dirty_disks.is_empty() { let mut dirty_disks = self.dirty_disks.write().await; for disk in &update.clear_dirty_disks { - dirty_disks.remove(disk); + // Only clear marks that predate the scan behind this update: a + // write landing while the walker was already past its prefix + // re-marks the disk, and that mark must survive the commit or + // the next dirty-subset refresh serves stale bytes as exact + // (backlog#1020 S19). + let clearable = match (update.scan_started_at, dirty_disks.get(disk)) { + (Some(scan_start), Some(marked_at)) => *marked_at <= scan_start, + _ => true, + }; + if clearable { + dirty_disks.remove(disk); + } } record_capacity_dirty_disk_count(dirty_disks.len()); } @@ -815,8 +836,9 @@ impl HybridCapacityManager { return; } + let now = Instant::now(); let mut dirty_disks = self.dirty_disks.write().await; - dirty_disks.extend(scope.disks.iter().cloned()); + dirty_disks.extend(scope.disks.iter().cloned().map(|disk| (disk, now))); record_capacity_dirty_disk_count(dirty_disks.len()); } @@ -908,7 +930,22 @@ impl HybridCapacityManager { pub async fn get_dirty_disks(&self) -> Vec { self.sync_global_dirty_scopes().await; let dirty_disks = self.dirty_disks.read().await; - dirty_disks.iter().cloned().collect() + dirty_disks.keys().cloned().collect() + } + + /// Drop dirty marks for disks outside the current local topology. + /// + /// The write side records every disk of an EC set — including remote + /// peers — while scans and dirty clearing only ever cover local disks, so + /// remote or removed disks would otherwise stay marked forever and keep + /// the dirty-disk gauge permanently non-zero (backlog#1020 S30). + pub async fn retain_dirty_disks_within(&self, local: &HashSet) { + let mut dirty_disks = self.dirty_disks.write().await; + let before = dirty_disks.len(); + dirty_disks.retain(|disk, _| local.contains(disk)); + if dirty_disks.len() != before { + record_capacity_dirty_disk_count(dirty_disks.len()); + } } /// Returns true if the manager has a complete per-disk cache and can safely refresh only dirty disks. @@ -1773,6 +1810,7 @@ mod tests { expected_disk_count: Some(2), replaces_disk_cache: true, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::RealTime, ) @@ -1797,6 +1835,7 @@ mod tests { expected_disk_count: Some(1), replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::WriteTriggered, ) @@ -1840,6 +1879,7 @@ mod tests { expected_disk_count: Some(2), replaces_disk_cache: true, clear_dirty_disks: Vec::new(), + scan_started_at: None, } } @@ -1864,6 +1904,7 @@ mod tests { expected_disk_count: None, replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::Scheduled, ) @@ -1906,6 +1947,7 @@ mod tests { expected_disk_count: None, replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::Scheduled, ) @@ -1934,6 +1976,7 @@ mod tests { expected_disk_count: None, replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::Scheduled, ) @@ -1947,6 +1990,84 @@ mod tests { assert!(!manager.can_refresh_dirty_subset().await); } + fn scope_disk(endpoint: &str, drive_path: &str) -> CapacityScopeDisk { + CapacityScopeDisk { + endpoint: endpoint.to_string(), + drive_path: drive_path.to_string(), + } + } + + #[tokio::test] + #[serial] + async fn test_commit_keeps_dirty_marks_recorded_after_scan_start() { + let manager = create_isolated_manager(HybridStrategyConfig::default()); + let disk = scope_disk("node-a", "/tmp/disk-a"); + + // Dirty before the scan starts: this mark is what the scan observed. + manager + .mark_dirty_scope(&CapacityScope { + disks: vec![disk.clone()], + }) + .await; + let scan_started_at = Instant::now(); + + // A write lands while the scan is in flight and re-marks the disk; + // the commit must not clear this newer mark. + tokio::time::sleep(Duration::from_millis(5)).await; + manager + .mark_dirty_scope(&CapacityScope { + disks: vec![disk.clone()], + }) + .await; + + let mut update = CapacityUpdate::exact(100, 1); + update.clear_dirty_disks = vec![disk.clone()]; + update.scan_started_at = Some(scan_started_at); + manager.update_capacity(update, DataSource::Scheduled).await; + + assert_eq!(manager.get_dirty_disks().await, vec![disk]); + } + + #[tokio::test] + #[serial] + async fn test_commit_clears_dirty_marks_recorded_before_scan_start() { + let manager = create_isolated_manager(HybridStrategyConfig::default()); + let disk = scope_disk("node-a", "/tmp/disk-a"); + + manager + .mark_dirty_scope(&CapacityScope { + disks: vec![disk.clone()], + }) + .await; + tokio::time::sleep(Duration::from_millis(5)).await; + + let mut update = CapacityUpdate::exact(100, 1); + update.clear_dirty_disks = vec![disk.clone()]; + update.scan_started_at = Some(Instant::now()); + manager.update_capacity(update, DataSource::Scheduled).await; + + assert!(manager.get_dirty_disks().await.is_empty()); + } + + #[tokio::test] + #[serial] + async fn test_retain_dirty_disks_within_drops_ghost_entries() { + let manager = create_isolated_manager(HybridStrategyConfig::default()); + let local = scope_disk("node-a", "/tmp/disk-a"); + let remote = scope_disk("http://peer:9000", "/data/disk-1"); + + manager + .mark_dirty_scope(&CapacityScope { + disks: vec![local.clone(), remote], + }) + .await; + + let local_set: HashSet = [local.clone()].into_iter().collect(); + manager.retain_dirty_disks_within(&local_set).await; + + assert_eq!(manager.get_dirty_disks().await, vec![local]); + } + #[tokio::test] #[serial] async fn test_spawn_refresh_recovers_from_construction_panic() { @@ -2078,6 +2199,7 @@ mod tests { expected_disk_count: Some(2), replaces_disk_cache: true, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::RealTime, ) @@ -2101,6 +2223,7 @@ mod tests { expected_disk_count: Some(1), replaces_disk_cache: false, clear_dirty_disks: Vec::new(), + scan_started_at: None, }; let returned = manager diff --git a/crates/object-capacity/src/scan.rs b/crates/object-capacity/src/scan.rs index f9f383fbf..8829c7633 100644 --- a/crates/object-capacity/src/scan.rs +++ b/crates/object-capacity/src/scan.rs @@ -312,6 +312,12 @@ pub async fn select_capacity_refresh_disks( capacity_manager: &HybridCapacityManager, disks: &[CapacityDiskRef], ) -> (Vec, bool) { + // The write side marks every disk of an EC set dirty, including remote + // peers, but only local disks are ever scanned and cleared — drop ghost + // entries so the dirty gauge reflects local pending work (backlog#1020). + let local_set: HashSet = disks.iter().map(disk_scope_key).collect(); + capacity_manager.retain_dirty_disks_within(&local_set).await; + if !capacity_manager.can_refresh_dirty_subset().await { return (disks.to_vec(), false); } @@ -336,6 +342,7 @@ pub async fn select_capacity_refresh_disks( } pub async fn refresh_capacity_with_scope(disks: Vec, dirty_subset: bool) -> Result { + let scan_started_at = Instant::now(); let report = calculate_data_dir_used_capacity_report(&disks) .await .map_err(|e| e.to_string())?; @@ -344,7 +351,9 @@ pub async fn refresh_capacity_with_scope(disks: Vec, dirty_subs return Err("dirty subset refresh had partial errors".to_string()); } - Ok(report.into_capacity_update(disks.len(), !dirty_subset)) + let mut update = report.into_capacity_update(disks.len(), !dirty_subset); + update.scan_started_at = Some(scan_started_at); + Ok(update) } /// Scan the provided local disk roots and return a summarized used-capacity result. @@ -1564,6 +1573,7 @@ mod tests { expected_disk_count: Some(2), replaces_disk_cache: true, clear_dirty_disks: Vec::new(), + scan_started_at: None, }, DataSource::RealTime, )