From 73a30178f506fa38ce75f5a37c35b07acb829174 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Thu, 9 Jul 2026 00:53:35 +0800 Subject: [PATCH] fix(object-capacity): make timeout fallback keep accumulated work and enter sampling within the time budget (#4517) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The capacity scan's timeout fallback required the exact prefix (threshold + sample_rate files, default ~200k) to be fully enumerated before any sample existed. Cold-cache HDDs cannot stat that many files inside the 15s budget, so large slow disks — exactly the disks sampling was designed for — always timed out with sampled_count == 0, discarded the accumulated exact-prefix work with a hard error, and re-ran 15s of useless full-disk I/O every 120s (S03). When a sample did exist, the estimate only extrapolated over files the walker had seen, silently under-reporting disks whose tail was never reached (S10). - Enter sampling early: once half the (possibly dynamic) time budget is consumed and the exact prefix hasn't filled, freeze the threshold at the current position so the remaining budget collects samples and the timeout fallback is always reachable. - Compensate the unseen tail: when the scan root is a dedicated mount point (stat st_dev differs from the parent), take the max of the seen-files extrapolation and the filesystem-level used bytes. On shared filesystems (multi-disk dev layouts) statvfs would overcount, so it is not trusted and behavior is unchanged. - sampled_count == 0 at timeout no longer hard-fails when a dedicated mount provides filesystem usage; the exact prefix is kept as a floor. Without any estimator the original error still surfaces. - Extract ScanLimits from env reading so the blocking scanner is parameterizable in tests. Ref: rustfs/backlog#1013 (S03+S10 from audit rustfs/backlog#1010) --- crates/object-capacity/Cargo.toml | 2 +- crates/object-capacity/src/scan.rs | 405 ++++++++++++++++++++++------- 2 files changed, 308 insertions(+), 99 deletions(-) diff --git a/crates/object-capacity/Cargo.toml b/crates/object-capacity/Cargo.toml index b48d4cadf..3fa0c616a 100644 --- a/crates/object-capacity/Cargo.toml +++ b/crates/object-capacity/Cargo.toml @@ -37,7 +37,7 @@ workspace = true [dependencies] rustfs-config = { workspace = true, features = ["constants"] } rustfs-io-metrics = { workspace = true } -rustfs-utils = { workspace = true } +rustfs-utils = { workspace = true, features = ["os"] } futures = { workspace = true } tokio = { workspace = true, features = ["sync", "time"] } tracing = { workspace = true } diff --git a/crates/object-capacity/src/scan.rs b/crates/object-capacity/src/scan.rs index 75ad85f11..ef0760a0c 100644 --- a/crates/object-capacity/src/scan.rs +++ b/crates/object-capacity/src/scan.rs @@ -547,41 +547,176 @@ impl ProgressMonitor { Ok(()) } - fn record_timeout_fallback(&self) { - record_capacity_timeout_fallback(); + /// True once at least half of the current (possibly dynamic) time budget + /// has elapsed — the signal to stop growing the exact prefix and start + /// sampling, so that when the timeout fires a usable sample exists even on + /// disks too slow to enumerate the full exact prefix (backlog#1013 S03). + fn half_budget_elapsed(&mut self, file_count: usize, avg_file_size: u64) -> bool { + let budget = if self.enable_dynamic_timeout { + self.calculate_dynamic_timeout(file_count, avg_file_size) + } else { + self.timeout + }; + self.start_time.elapsed() >= budget / 2 } } +/// Scan configuration resolved from the environment, separated so the blocking +/// scanner can be exercised with explicit limits in tests. +#[derive(Debug, Clone)] +struct ScanLimits { + max_files_threshold: usize, + base_timeout: Duration, + min_timeout: Duration, + max_timeout: Duration, + stall_timeout: Duration, + sample_rate: usize, + enable_dynamic_timeout: bool, + follow_symlinks: bool, + max_symlink_depth: u8, +} + +impl ScanLimits { + fn from_env() -> Self { + let sample_rate = get_sample_rate(); + let effective_sample_rate = if sample_rate == 0 { + warn!( + event = EVENT_CAPACITY_SCAN_SAMPLING_CLAMPED, + component = LOG_COMPONENT_CAPACITY, + subsystem = LOG_SUBSYSTEM_SAMPLING, + result = "clamped", + configured_sample_rate = 0, + effective_sample_rate = 1, + reason = "zero_sample_rate", + "capacity scan sampling configuration clamped" + ); + 1 + } else { + sample_rate + }; + + Self { + max_files_threshold: get_max_files_threshold(), + base_timeout: get_stat_timeout(), + min_timeout: get_min_timeout(), + max_timeout: get_max_timeout(), + stall_timeout: get_stall_timeout(), + sample_rate: effective_sample_rate, + enable_dynamic_timeout: get_enable_dynamic_timeout(), + follow_symlinks: get_follow_symlinks(), + max_symlink_depth: get_max_symlink_depth(), + } + } +} + +/// Filesystem-level used bytes for `path`, only when `path` is a dedicated +/// mount point. On a shared filesystem (several logical disks per filesystem, +/// common in dev deployments) statvfs counts unrelated data and would badly +/// overcount, so `None` is returned and callers keep scan-based estimates. +fn mount_point_used_bytes(path: &Path) -> Option { + let parent = path.parent()?; + let dedicated = !rustfs_utils::os::same_disk(path.to_str()?, parent.to_str()?).ok()?; + if !dedicated { + return None; + } + rustfs_utils::os::get_info(path).ok().map(|info| info.used) +} + +/// Pick the timeout-fallback value from the two independent estimators. +/// +/// The seen-files extrapolation guarantees at least the observed data; the +/// filesystem-level usage covers the walker's unseen tail. When both exist the +/// max wins (both are lower..upper bounds of the same truth from opposite +/// sides). With only filesystem usage, the exact prefix is still a hard floor. +/// `None` means no trustworthy estimator exists. +fn combine_timeout_estimates(seen_estimate: Option, disk_used: Option, exact_prefix_bytes: u64) -> Option { + match (seen_estimate, disk_used) { + (Some(seen), Some(disk)) => Some(seen.max(disk)), + (Some(seen), None) => Some(seen), + (None, Some(disk)) => Some(disk.max(exact_prefix_bytes)), + (None, None) => None, + } +} + +/// Build the best available estimate when a scan times out or stalls, instead +/// of discarding the accumulated work with a hard error (backlog#1013 S03+S10). +/// +/// The seen-files extrapolation only covers files the walker reached; the +/// filesystem-level usage covers the whole disk but is only trustworthy on a +/// dedicated mount. Taking the max compensates for the walker's unseen tail +/// without regressing shared-filesystem deployments. Returns `None` when +/// neither source is available (no sample and no dedicated mount) — the caller +/// then propagates the original error, as before. +#[allow(clippy::too_many_arguments)] +fn timeout_fallback_estimate( + path: &Path, + reason: &'static str, + exact_prefix_bytes: u64, + overflow_sampled_bytes: u64, + file_count: usize, + sampled_count: usize, + effective_threshold: usize, + had_partial_errors: bool, + start_time: Instant, +) -> Option { + let seen_estimate = (sampled_count > 0).then(|| { + let overflow_count = file_count.saturating_sub(effective_threshold); + let estimated_overflow = estimate_overflow_bytes(overflow_sampled_bytes, overflow_count as u64, sampled_count as u64); + exact_prefix_bytes.saturating_add(estimated_overflow) + }); + let disk_used = mount_point_used_bytes(path); + let used_bytes = combine_timeout_estimates(seen_estimate, disk_used, exact_prefix_bytes)?; + + info!( + event = EVENT_CAPACITY_SCAN_SAMPLING_APPLIED, + component = LOG_COMPONENT_CAPACITY, + subsystem = LOG_SUBSYSTEM_SAMPLING, + result = "fallback_estimate", + reason, + file_count, + exact_prefix_bytes, + sampled_count, + seen_estimate_bytes = seen_estimate, + disk_used_bytes = disk_used, + estimated_total_bytes = used_bytes, + "capacity scan sampling applied" + ); + record_capacity_timeout_fallback(); + record_capacity_scan_sampling(sampled_count, true); + record_capacity_scan_mode("timeout_fallback"); + + Some(CapacityScanResult { + used_bytes, + file_count, + sampled_count, + is_estimated: true, + scan_duration: start_time.elapsed(), + had_partial_errors, + }) +} + async fn get_dir_size_async(path: &Path) -> Result { let path = path.to_path_buf(); + let limits = ScanLimits::from_env(); - let max_files_threshold = get_max_files_threshold(); - let base_timeout = get_stat_timeout(); - let min_timeout = get_min_timeout(); - let max_timeout = get_max_timeout(); - let stall_timeout = get_stall_timeout(); - let sample_rate = get_sample_rate(); - let enable_dynamic_timeout = get_enable_dynamic_timeout(); - let follow_symlinks = get_follow_symlinks(); - let max_symlink_depth = get_max_symlink_depth(); + tokio::task::spawn_blocking(move || scan_dir_blocking(&path, &limits)) + .await + .map_err(std::io::Error::other)? +} - let effective_sample_rate = if sample_rate == 0 { - warn!( - event = EVENT_CAPACITY_SCAN_SAMPLING_CLAMPED, - component = LOG_COMPONENT_CAPACITY, - subsystem = LOG_SUBSYSTEM_SAMPLING, - result = "clamped", - configured_sample_rate = 0, - effective_sample_rate = 1, - reason = "zero_sample_rate", - "capacity scan sampling configuration clamped" - ); - 1 - } else { - sample_rate - }; - - tokio::task::spawn_blocking(move || { +fn scan_dir_blocking(path: &Path, limits: &ScanLimits) -> Result { + let ScanLimits { + max_files_threshold, + base_timeout, + min_timeout, + max_timeout, + stall_timeout, + sample_rate: effective_sample_rate, + enable_dynamic_timeout, + follow_symlinks, + max_symlink_depth, + } = *limits; + { if !path.exists() { return Err(std::io::Error::new( std::io::ErrorKind::NotFound, @@ -596,12 +731,15 @@ async fn get_dir_size_async(path: &Path) -> Result Result meta, Err(err) => { - if is_tmp_trash_metadata_not_found(&path, entry.path(), err.io_error().map(|err| err.kind())) { + if is_tmp_trash_metadata_not_found(path, entry.path(), err.io_error().map(|err| err.kind())) { debug!( event = EVENT_CAPACITY_SCAN_METADATA_FAILED, component = LOG_COMPONENT_CAPACITY, @@ -675,7 +813,7 @@ async fn get_dir_size_async(path: &Path) -> Result 0 { exact_prefix_bytes / exact_count as u64 } else { @@ -685,47 +823,54 @@ async fn get_dir_size_async(path: &Path) -> Result= CAPACITY_PROGRESS_CHECK_STRIDE; - if should_check_progress && let Err(e) = progress_monitor.update_and_check_timeout(file_count, avg_size) { - if sampled_count > 0 { - let overflow_count = file_count.saturating_sub(max_files_threshold); - let estimated_overflow = - estimate_overflow_bytes(overflow_sampled_bytes, overflow_count as u64, sampled_count as u64); - let estimated_total = exact_prefix_bytes.saturating_add(estimated_overflow); + if should_check_progress { + if let Err(e) = progress_monitor.update_and_check_timeout(file_count, avg_size) { + if let Some(result) = timeout_fallback_estimate( + path, + "timeout_or_stall", + exact_prefix_bytes, + overflow_sampled_bytes, + file_count, + sampled_count, + effective_threshold, + had_partial_errors, + start_time, + ) { + return Ok(result); + } + return Err(e); + } + + // Half the time budget is gone and the exact prefix hasn't + // filled: freeze the threshold at the current position so the + // remaining budget collects overflow samples. Without this, a + // disk too slow to enumerate the full prefix within the budget + // reaches the timeout with sampled_count == 0 and used to lose + // the whole scan (backlog#1013 S03). + if file_count < effective_threshold && progress_monitor.half_budget_elapsed(file_count, avg_size) { + effective_threshold = file_count; info!( event = EVENT_CAPACITY_SCAN_SAMPLING_APPLIED, component = LOG_COMPONENT_CAPACITY, subsystem = LOG_SUBSYSTEM_SAMPLING, - result = "fallback_estimate", - reason = "timeout_or_stall", + result = "early_sampling", + reason = "time_budget", file_count, exact_prefix_bytes, - estimated_overflow_bytes = estimated_overflow, - sampled_count, - estimated_total_bytes = estimated_total, + frozen_threshold = effective_threshold, + configured_threshold = max_files_threshold, "capacity scan sampling applied" ); - progress_monitor.record_timeout_fallback(); - record_capacity_scan_sampling(sampled_count, true); - record_capacity_scan_mode("timeout_fallback"); - return Ok(CapacityScanResult { - used_bytes: estimated_total, - file_count, - sampled_count, - is_estimated: true, - scan_duration: start_time.elapsed(), - had_partial_errors, - }); + record_capacity_scan_mode("early_sampling"); } - return Err(e); - } - if should_check_progress { + last_progress_check_files = file_count; } - if file_count <= max_files_threshold { + if file_count <= effective_threshold { exact_prefix_bytes += metadata.len(); } else { - let overflow_index = file_count - max_files_threshold; + let overflow_index = file_count - effective_threshold; if overflow_index.is_multiple_of(effective_sample_rate) { overflow_sampled_bytes += metadata.len(); sampled_count += 1; @@ -748,7 +893,7 @@ async fn get_dir_size_async(path: &Path) -> Result last_progress_check_files { - let exact_count = file_count.min(max_files_threshold); + let exact_count = file_count.min(effective_threshold); let avg_size = if exact_count > 0 { exact_prefix_bytes / exact_count as u64 } else { @@ -756,35 +901,18 @@ async fn get_dir_size_async(path: &Path) -> Result 0 { - let overflow_count = file_count.saturating_sub(max_files_threshold); - let estimated_overflow = - estimate_overflow_bytes(overflow_sampled_bytes, overflow_count as u64, sampled_count as u64); - let estimated_total = exact_prefix_bytes.saturating_add(estimated_overflow); - info!( - event = EVENT_CAPACITY_SCAN_SAMPLING_APPLIED, - component = LOG_COMPONENT_CAPACITY, - subsystem = LOG_SUBSYSTEM_SAMPLING, - result = "fallback_estimate", - reason = "final_timeout_or_stall", - file_count, - exact_prefix_bytes, - estimated_overflow_bytes = estimated_overflow, - sampled_count, - estimated_total_bytes = estimated_total, - "capacity scan sampling applied" - ); - progress_monitor.record_timeout_fallback(); - record_capacity_scan_sampling(sampled_count, true); - record_capacity_scan_mode("timeout_fallback"); - return Ok(CapacityScanResult { - used_bytes: estimated_total, - file_count, - sampled_count, - is_estimated: true, - scan_duration: start_time.elapsed(), - had_partial_errors, - }); + if let Some(result) = timeout_fallback_estimate( + path, + "final_timeout_or_stall", + exact_prefix_bytes, + overflow_sampled_bytes, + file_count, + sampled_count, + effective_threshold, + had_partial_errors, + start_time, + ) { + return Ok(result); } return Err(e); } @@ -803,8 +931,8 @@ async fn get_dir_size_async(path: &Path) -> Result max_files_threshold && sampled_count > 0 { - let overflow_count = file_count - max_files_threshold; + if file_count > effective_threshold && sampled_count > 0 { + let overflow_count = file_count - effective_threshold; let estimated_overflow = estimate_overflow_bytes(overflow_sampled_bytes, overflow_count as u64, sampled_count as u64); let estimated_size = exact_prefix_bytes.saturating_add(estimated_overflow); info!( @@ -814,7 +942,7 @@ async fn get_dir_size_async(path: &Path) -> Result Result max_files_threshold { - let overflow_count = file_count - max_files_threshold; - let exact_prefix_count = file_count.min(max_files_threshold) as u64; + } else if file_count > effective_threshold { + let overflow_count = file_count - effective_threshold; + let exact_prefix_count = file_count.min(effective_threshold) as u64; let avg_prefix_size = exact_prefix_bytes.checked_div(exact_prefix_count).unwrap_or(0); let estimated_overflow = avg_prefix_size.saturating_mul(overflow_count as u64); let estimated_size = exact_prefix_bytes.saturating_add(estimated_overflow); @@ -845,7 +973,7 @@ async fn get_dir_size_async(path: &Path) -> Result Result ScanLimits { + ScanLimits { + max_files_threshold: 100_000, + base_timeout, + min_timeout: Duration::from_millis(0), + max_timeout: base_timeout, + stall_timeout: Duration::from_secs(600), + sample_rate: 1, + enable_dynamic_timeout: false, + follow_symlinks: false, + max_symlink_depth: 40, + } + } + + #[test] + fn test_scan_dir_blocking_zero_budget_without_estimators_still_errors() { + use std::fs::File; + use std::io::Write; + + let temp_dir = tempfile::TempDir::new().unwrap(); + for i in 0..3 { + let mut file = File::create(temp_dir.path().join(format!("f{i}"))).unwrap(); + file.write_all(b"x").unwrap(); + } + + // Zero budget times out at the first progress check with no samples; + // a tempdir is not a dedicated mount, so no estimator exists and the + // original timeout error must still surface (last-resort behavior). + let result = scan_dir_blocking(temp_dir.path(), &tight_limits(Duration::ZERO)); + assert_eq!(result.unwrap_err().kind(), std::io::ErrorKind::TimedOut); + } + + #[test] + fn test_scan_dir_blocking_generous_budget_still_exact() { + use std::fs::File; + use std::io::Write; + + let temp_dir = tempfile::TempDir::new().unwrap(); + for i in 0..5 { + let mut file = File::create(temp_dir.path().join(format!("f{i}"))).unwrap(); + file.write_all(b"hello").unwrap(); + } + + let result = scan_dir_blocking(temp_dir.path(), &tight_limits(Duration::from_secs(600))).unwrap(); + assert_eq!(result.used_bytes, 25); + assert_eq!(result.file_count, 5); + assert!(!result.is_estimated); + assert!(!result.had_partial_errors); + } + #[tokio::test] async fn test_partial_scan_report_produces_degraded_update() { use std::fs::File;