fix(capacity): back off timed-out scans (#5770)

* fix(capacity): back off timed-out scans

* fix(capacity): guard incomplete metadata baselines

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: Anthony Martin <949506+anthonymartin@users.noreply.github.com>
Co-authored-by: houseme <housemecn@gmail.com>
Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
anthonymartin
2026-08-06 08:56:45 -07:00
committed by GitHub
parent 5e7e25b7d1
commit 77f2b948c2
3 changed files with 361 additions and 36 deletions
+153 -26
View File
@@ -55,7 +55,7 @@ const EVENT_CAPACITY_REFRESH_CANCELLED: &str = "capacity_refresh_cancelled";
const EVENT_CAPACITY_REFRESH_RUNTIME_SUMMARY: &str = "capacity_refresh_runtime_summary";
const EVENT_CAPACITY_REFRESH_INTERVAL_CLAMPED: &str = "capacity_refresh_interval_clamped";
const EVENT_CAPACITY_REFRESH_SCHEDULED: &str = "capacity_refresh_scheduled";
const EVENT_CAPACITY_REFRESH_SKIPPED: &str = "capacity_refresh_skipped";
const MAX_SCHEDULED_REFRESH_BACKOFF: Duration = Duration::from_secs(30 * 60);
// ============================================================================
// Configuration Functions
@@ -335,6 +335,8 @@ pub struct CapacityUpdate {
/// `is_estimated` only reflects sampling; this flag is the only carrier of
/// the "some disks failed to scan" fact (backlog#1014).
pub degraded: bool,
/// Whether the scan reached its time budget and returned a fallback estimate.
pub timed_out: bool,
/// Per-disk breakdown captured from a successful refresh.
pub per_disk: Vec<DiskCapacityUpdate>,
/// Expected disk count for a complete disk cache.
@@ -357,6 +359,7 @@ impl CapacityUpdate {
file_count,
is_estimated: false,
degraded: false,
timed_out: false,
per_disk: Vec::new(),
expected_disk_count: None,
replaces_disk_cache: false,
@@ -372,6 +375,7 @@ impl CapacityUpdate {
file_count,
is_estimated: true,
degraded: false,
timed_out: false,
per_disk: Vec::new(),
expected_disk_count: None,
replaces_disk_cache: false,
@@ -387,6 +391,7 @@ impl CapacityUpdate {
file_count: 0,
is_estimated: true,
degraded: false,
timed_out: false,
per_disk: Vec::new(),
expected_disk_count: None,
replaces_disk_cache: false,
@@ -1321,6 +1326,51 @@ fn clamp_background_interval(value: Duration, env_var: &'static str) -> Duration
clamped
}
#[derive(Debug)]
struct ScheduledRefreshBackoff {
base: Duration,
current: Duration,
max: Duration,
}
impl ScheduledRefreshBackoff {
fn new(base: Duration) -> Self {
Self {
base,
current: base,
max: base.max(MAX_SCHEDULED_REFRESH_BACKOFF),
}
}
fn delay(&self) -> Duration {
self.current
}
fn record_result(&mut self, clean: bool) {
self.current = if clean {
self.base
} else {
self.current.saturating_mul(2).min(self.max)
};
}
}
fn scheduled_refresh_was_clean(result: &Result<CapacityUpdate, String>) -> bool {
matches!(result, Ok(update) if !update.timed_out && !update.degraded)
}
async fn run_scheduled_refresh_loop<F, Fut>(refresh_interval: Duration, mut refresh: F)
where
F: FnMut() -> Fut,
Fut: Future<Output = bool>,
{
let mut backoff = ScheduledRefreshBackoff::new(refresh_interval);
loop {
tokio::time::sleep(backoff.delay()).await;
backoff.record_result(refresh().await);
}
}
pub async fn start_background_task(disks: Vec<CapacityDiskRef>) {
let manager = get_capacity_manager();
let manager_for_refresh = manager.clone();
@@ -1332,24 +1382,12 @@ pub async fn start_background_task(disks: Vec<CapacityDiskRef>) {
metrics_interval = clamp_background_interval(metrics_interval, ENV_CAPACITY_METRICS_INTERVAL);
tokio::spawn(async move {
let mut timer = tokio::time::interval_at(tokio::time::Instant::now() + refresh_interval, refresh_interval);
loop {
timer.tick().await;
run_scheduled_refresh_loop(refresh_interval, move || {
let start = Instant::now();
let manager = manager_for_refresh.clone();
let disks = disks.clone();
let disk_count = disks.len();
let started = manager
.clone()
.spawn_refresh_if_needed(
DataSource::Scheduled,
move || async move { refresh_capacity_with_scope(disks, false).await },
)
.await;
if started {
async move {
debug!(
event = EVENT_CAPACITY_REFRESH_SCHEDULED,
component = LOG_COMPONENT_CAPACITY,
@@ -1360,18 +1398,16 @@ pub async fn start_background_task(disks: Vec<CapacityDiskRef>) {
enqueue_latency_ms = start.elapsed().as_millis() as u64,
"capacity refresh scheduled"
);
} else {
debug!(
event = EVENT_CAPACITY_REFRESH_SKIPPED,
component = LOG_COMPONENT_CAPACITY,
subsystem = LOG_SUBSYSTEM_RUNTIME,
state = "inflight",
source = DataSource::Scheduled.as_metric_label(),
disk_count,
"capacity refresh skipped"
);
let result = manager
.refresh_or_join(
DataSource::Scheduled,
move || async move { refresh_capacity_with_scope(disks, false).await },
)
.await;
scheduled_refresh_was_clean(&result)
}
}
})
.await;
});
tokio::spawn(async move {
@@ -1400,6 +1436,89 @@ mod tests {
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
#[test]
fn test_scheduled_refresh_backoff_grows_caps_and_resets() {
let mut backoff = ScheduledRefreshBackoff::new(Duration::from_secs(120));
let cases = [
(false, 240),
(false, 480),
(false, 960),
(false, 1_800),
(false, 1_800),
(true, 120),
];
assert_eq!(backoff.delay(), Duration::from_secs(120));
for (clean, expected_secs) in cases {
backoff.record_result(clean);
assert_eq!(backoff.delay(), Duration::from_secs(expected_secs));
}
let mut configured_above_cap = ScheduledRefreshBackoff::new(Duration::from_secs(3_600));
configured_above_cap.record_result(false);
assert_eq!(configured_above_cap.delay(), Duration::from_secs(3_600));
}
#[test]
fn test_scheduled_refresh_cleanliness_distinguishes_sampling_from_timeout() {
let exact = Ok(CapacityUpdate::exact(10, 1));
let sampled = Ok(CapacityUpdate::estimated(10, 1));
let mut timed_out = CapacityUpdate::estimated(10, 1);
timed_out.timed_out = true;
let mut degraded = CapacityUpdate::exact(10, 1);
degraded.degraded = true;
assert!(scheduled_refresh_was_clean(&exact));
assert!(scheduled_refresh_was_clean(&sampled));
assert!(!scheduled_refresh_was_clean(&Ok(timed_out)));
assert!(!scheduled_refresh_was_clean(&Ok(degraded)));
assert!(!scheduled_refresh_was_clean(&Err("scan failed".to_string())));
}
#[tokio::test(start_paused = true)]
async fn test_scheduled_refresh_loop_applies_backoff_and_reset() {
use std::collections::VecDeque;
use std::sync::Mutex as StdMutex;
let calls = Arc::new(AtomicUsize::new(0));
let results = Arc::new(StdMutex::new(VecDeque::from([false, false, true, true])));
let task_calls = calls.clone();
let task_results = results.clone();
let task = tokio::spawn(run_scheduled_refresh_loop(Duration::from_secs(10), move || {
let task_calls = task_calls.clone();
let task_results = task_results.clone();
async move {
task_calls.fetch_add(1, Ordering::SeqCst);
task_results
.lock()
.expect("scheduled result queue lock should succeed")
.pop_front()
.unwrap_or(true)
}
}));
tokio::task::yield_now().await;
tokio::time::advance(Duration::from_secs(10)).await;
tokio::task::yield_now().await;
assert_eq!(calls.load(Ordering::SeqCst), 1);
tokio::time::advance(Duration::from_secs(19)).await;
tokio::task::yield_now().await;
assert_eq!(calls.load(Ordering::SeqCst), 1);
tokio::time::advance(Duration::from_secs(1)).await;
tokio::task::yield_now().await;
assert_eq!(calls.load(Ordering::SeqCst), 2);
tokio::time::advance(Duration::from_secs(40)).await;
tokio::task::yield_now().await;
assert_eq!(calls.load(Ordering::SeqCst), 3);
tokio::time::advance(Duration::from_secs(10)).await;
tokio::task::yield_now().await;
assert_eq!(calls.load(Ordering::SeqCst), 4);
task.abort();
}
type ConfigGetterCase = (&'static str, fn() -> u64, u64, &'static str, u64);
/// Table of env-configurable getters: (env var, getter normalized to u64,
@@ -1926,6 +2045,7 @@ mod tests {
file_count: 3,
is_estimated: false,
degraded: false,
timed_out: false,
per_disk: vec![
DiskCapacityUpdate {
disk: CapacityScopeDisk {
@@ -1962,6 +2082,7 @@ mod tests {
file_count: 1,
is_estimated: true,
degraded: false,
timed_out: false,
per_disk: vec![DiskCapacityUpdate {
disk: CapacityScopeDisk {
endpoint: "node-a".to_string(),
@@ -2011,6 +2132,7 @@ mod tests {
file_count: 2,
is_estimated: false,
degraded: false,
timed_out: false,
per_disk: vec![
disk_entry("node-a", "/tmp/disk-a", 100, 1),
disk_entry("node-b", "/tmp/disk-b", 100, 1),
@@ -2039,6 +2161,7 @@ mod tests {
file_count: 1,
is_estimated: false,
degraded: true,
timed_out: false,
per_disk: vec![disk_entry("node-a", "/tmp/disk-a", 100, 1)],
expected_disk_count: None,
replaces_disk_cache: false,
@@ -2082,6 +2205,7 @@ mod tests {
file_count: 1,
is_estimated: false,
degraded: true,
timed_out: false,
per_disk: Vec::new(),
expected_disk_count: None,
replaces_disk_cache: false,
@@ -2111,6 +2235,7 @@ mod tests {
file_count: 1,
is_estimated: false,
degraded: true,
timed_out: false,
per_disk: vec![disk_entry("node-a", "/tmp/disk-a", 100, 1)],
expected_disk_count: None,
replaces_disk_cache: false,
@@ -2315,6 +2440,7 @@ mod tests {
file_count: 3,
is_estimated: false,
degraded: false,
timed_out: false,
per_disk: vec![
DiskCapacityUpdate {
disk: CapacityScopeDisk {
@@ -2350,6 +2476,7 @@ mod tests {
file_count: 1,
is_estimated: true,
degraded: false,
timed_out: false,
per_disk: vec![DiskCapacityUpdate {
disk: CapacityScopeDisk {
endpoint: "node-a".to_string(),
+204 -10
View File
@@ -92,9 +92,13 @@ impl CapacityScanReport {
is_estimated: entry.scan.is_estimated,
})
.collect();
update.expected_disk_count = Some(expected_disk_count);
update.replaces_disk_cache = replaces_disk_cache;
update.clear_dirty_disks = update.per_disk.iter().map(|entry| entry.disk.clone()).collect();
// Skipped metadata or timeout fallback estimates can update an
// existing complete cache, but must not establish a new baseline.
if !self.summary.timed_out && !self.summary.metadata_incomplete {
update.expected_disk_count = Some(expected_disk_count);
update.replaces_disk_cache = replaces_disk_cache;
update.clear_dirty_disks = update.per_disk.iter().map(|entry| entry.disk.clone()).collect();
}
} else if self.summary.had_partial_errors {
// Partial failure: mark the reading degraded and surface only the
// disks whose own scan fully succeeded, so the manager can merge
@@ -117,6 +121,8 @@ impl CapacityScanReport {
update.clear_dirty_disks = update.per_disk.iter().map(|entry| entry.disk.clone()).collect();
}
update.timed_out = self.summary.timed_out;
update
}
}
@@ -145,6 +151,10 @@ fn estimate_overflow_bytes(overflow_sampled_bytes: u64, overflow_count: u64, sam
scaled.min(u64::MAX as u128) as u64
}
fn overflow_entry_is_sampled(overflow_index: usize, sample_rate: usize) -> bool {
overflow_index.is_multiple_of(sample_rate)
}
fn disk_scope_key(disk: &CapacityDiskRef) -> CapacityScopeDisk {
CapacityScopeDisk {
endpoint: disk.endpoint.clone(),
@@ -197,6 +207,10 @@ async fn calculate_data_dir_used_capacity_report(
let mut has_failure = false;
let mut has_success = false;
let mut is_estimated = false;
let mut timed_out = false;
let mut metadata_incomplete = false;
#[cfg(test)]
let mut total_metadata_reads = 0usize;
let mut per_disk = Vec::with_capacity(disks.len());
let concurrency_limit = disks.len().clamp(1, MAX_CAPACITY_SCAN_CONCURRENCY);
@@ -232,6 +246,12 @@ async fn calculate_data_dir_used_capacity_report(
total_files += scan.file_count;
total_sampled += scan.sampled_count;
is_estimated |= scan.is_estimated;
timed_out |= scan.timed_out;
metadata_incomplete |= scan.metadata_incomplete;
#[cfg(test)]
{
total_metadata_reads += scan.metadata_reads;
}
has_failure |= scan.had_partial_errors;
has_success = true;
if let Some(disk) = disks
@@ -289,6 +309,10 @@ async fn calculate_data_dir_used_capacity_report(
is_estimated,
scan_duration: start.elapsed(),
had_partial_errors: false,
timed_out,
metadata_incomplete,
#[cfg(test)]
metadata_reads: total_metadata_reads,
};
if has_failure {
@@ -601,6 +625,10 @@ fn timeout_fallback_estimate(
is_estimated: true,
scan_duration: start_time.elapsed(),
had_partial_errors,
timed_out: true,
metadata_incomplete: false,
#[cfg(test)]
metadata_reads: 0,
})
}
@@ -681,6 +709,9 @@ fn scan_dir_blocking(path: &Path, limits: &ScanLimits, cancelled: &AtomicBool) -
let mut overflow_sampled_bytes = 0u64;
let mut file_count = 0usize;
let mut sampled_count = 0usize;
let mut metadata_incomplete = false;
#[cfg(test)]
let mut metadata_reads = 0usize;
let mut had_partial_errors = false;
let mut entries_visited = 0usize;
let mut last_progress_check_entries = 0usize;
@@ -810,6 +841,20 @@ fn scan_dir_blocking(path: &Path, limits: &ScanLimits, cancelled: &AtomicBool) -
continue;
}
let exact_entry = file_count < effective_threshold;
if !exact_entry {
file_count += 1;
let overflow_index = file_count - effective_threshold;
if !overflow_entry_is_sampled(overflow_index, effective_sample_rate) {
metadata_incomplete = true;
continue;
}
}
#[cfg(test)]
{
metadata_reads += 1;
}
let metadata = match entry.metadata() {
Ok(meta) => meta,
Err(err) => {
@@ -854,15 +899,12 @@ fn scan_dir_blocking(path: &Path, limits: &ScanLimits, cancelled: &AtomicBool) -
}
};
file_count += 1;
if file_count <= effective_threshold {
if exact_entry {
file_count += 1;
exact_prefix_bytes += metadata.len();
} else {
let overflow_index = file_count - effective_threshold;
if overflow_index.is_multiple_of(effective_sample_rate) {
overflow_sampled_bytes += metadata.len();
sampled_count += 1;
}
overflow_sampled_bytes += metadata.len();
sampled_count += 1;
if file_count.is_multiple_of(100_000) {
debug!(
@@ -934,6 +976,10 @@ fn scan_dir_blocking(path: &Path, limits: &ScanLimits, cancelled: &AtomicBool) -
is_estimated: true,
scan_duration: start_time.elapsed(),
had_partial_errors,
timed_out: false,
metadata_incomplete,
#[cfg(test)]
metadata_reads,
})
} else if file_count > effective_threshold {
let overflow_count = file_count - effective_threshold;
@@ -966,6 +1012,10 @@ fn scan_dir_blocking(path: &Path, limits: &ScanLimits, cancelled: &AtomicBool) -
is_estimated: true,
scan_duration: start_time.elapsed(),
had_partial_errors,
timed_out: false,
metadata_incomplete,
#[cfg(test)]
metadata_reads,
})
} else {
record_capacity_scan_sampling(0, false);
@@ -987,6 +1037,10 @@ fn scan_dir_blocking(path: &Path, limits: &ScanLimits, cancelled: &AtomicBool) -
is_estimated: false,
scan_duration: start_time.elapsed(),
had_partial_errors,
timed_out: false,
metadata_incomplete,
#[cfg(test)]
metadata_reads,
})
}
}
@@ -999,6 +1053,7 @@ mod tests {
use crate::capacity_scope::{CapacityScope, CapacityScopeDisk};
#[cfg(unix)]
use rustfs_config::ENV_CAPACITY_FOLLOW_SYMLINKS;
use rustfs_config::{ENV_CAPACITY_MAX_FILES_THRESHOLD, ENV_CAPACITY_SAMPLE_RATE};
use serial_test::serial;
/// Reference implementation using unbounded `u128` arithmetic, clamped to
@@ -1075,6 +1130,63 @@ mod tests {
assert_eq!(estimate_overflow_bytes(0, 10, 0), 0);
}
#[test]
fn test_overflow_sampling_avoids_unsampled_metadata_reads() {
let threshold = 200_000usize;
let file_count = 1_000_000usize;
let sample_rate = 200usize;
let metadata_reads = threshold
+ (1..=file_count - threshold)
.filter(|index| overflow_entry_is_sampled(*index, sample_rate))
.count();
assert_eq!(metadata_reads, 204_000);
assert!(metadata_reads < file_count / 4);
}
#[test]
fn test_sample_before_metadata_preserves_uniform_file_estimate() {
use std::fs::File;
use std::io::Write;
let temp_dir = tempfile::TempDir::new().expect("capacity scan tempdir should be created");
for index in 0..11 {
let mut file =
File::create(temp_dir.path().join(format!("f{index}"))).expect("capacity scan fixture file should be created");
file.write_all(b"12345")
.expect("capacity scan fixture bytes should be written");
}
let limits = ScanLimits {
max_files_threshold: 2,
base_timeout: Duration::from_secs(600),
min_timeout: Duration::from_secs(1),
max_timeout: Duration::from_secs(600),
sample_rate: 3,
enable_dynamic_timeout: false,
follow_symlinks: false,
};
let result =
scan_dir_blocking(temp_dir.path(), &limits, &AtomicBool::new(false)).expect("sampled capacity scan should succeed");
assert_eq!(result.file_count, 11);
assert_eq!(result.sampled_count, 3);
assert_eq!(result.metadata_reads, 5);
assert_eq!(result.used_bytes, 55);
assert!(result.is_estimated);
assert!(result.metadata_incomplete);
assert!(!result.timed_out);
}
#[test]
fn test_timeout_fallback_is_marked_for_scheduler_backoff() {
let temp_dir = tempfile::TempDir::new().expect("capacity scan tempdir should be created");
let result = timeout_fallback_estimate(temp_dir.path(), "timeout", 10, 5, 2, 1, 1, false, Instant::now())
.expect("sampled timeout should produce a fallback estimate");
assert!(result.timed_out);
assert!(result.is_estimated);
}
#[tokio::test]
async fn test_get_dir_size_async_empty_directory() {
use tempfile::TempDir;
@@ -1469,6 +1581,87 @@ mod tests {
assert!(!update.replaces_disk_cache);
}
#[test]
fn test_into_capacity_update_incomplete_results_do_not_replace_disk_cache() {
for scan in [
CapacityScanResult {
used_bytes: 100,
file_count: 10,
is_estimated: true,
metadata_incomplete: true,
..Default::default()
},
CapacityScanResult {
used_bytes: 100,
file_count: 10,
is_estimated: true,
timed_out: true,
..Default::default()
},
] {
let disk = CapacityScopeDisk {
endpoint: "node-a".to_string(),
drive_path: "/tmp/disk-a".to_string(),
};
let report = CapacityScanReport {
summary: scan,
per_disk: vec![DiskCapacityScanResult {
disk: disk.clone(),
scan,
}],
};
let update = report.into_capacity_update(1, true);
assert!(!update.degraded);
assert_eq!(update.per_disk.len(), 1);
assert_eq!(update.per_disk[0].disk, disk);
assert_eq!(update.expected_disk_count, None);
assert!(!update.replaces_disk_cache);
assert!(update.clear_dirty_disks.is_empty());
}
}
#[tokio::test]
#[serial]
async fn test_metadata_incomplete_aggregate_does_not_replace_disk_cache() {
use std::fs::File;
use std::io::Write;
let temp_dir = tempfile::TempDir::new().expect("capacity scan tempdir should be created");
for index in 0..11 {
let mut file =
File::create(temp_dir.path().join(format!("f{index}"))).expect("capacity scan fixture file should be created");
file.write_all(b"12345")
.expect("capacity scan fixture bytes should be written");
}
let disks = vec![CapacityDiskRef {
endpoint: "node-a".to_string(),
drive_path: temp_dir.path().display().to_string(),
}];
let report = temp_env::async_with_vars(
[
(ENV_CAPACITY_MAX_FILES_THRESHOLD, Some("2")),
(ENV_CAPACITY_SAMPLE_RATE, Some("3")),
],
calculate_data_dir_used_capacity_report(&disks),
)
.await
.expect("sampled capacity report should succeed");
assert!(report.summary.metadata_incomplete);
assert_eq!(report.per_disk.len(), 1);
assert!(report.per_disk[0].scan.metadata_incomplete);
let update = report.into_capacity_update(1, true);
assert_eq!(update.per_disk.len(), 1);
assert_eq!(update.expected_disk_count, None);
assert!(!update.replaces_disk_cache);
assert!(update.clear_dirty_disks.is_empty());
}
#[tokio::test]
async fn test_select_capacity_refresh_disks_returns_full_when_disk_cache_incomplete() {
let manager = create_isolated_manager(HybridStrategyConfig::default());
@@ -1507,6 +1700,7 @@ mod tests {
file_count: 3,
is_estimated: false,
degraded: false,
timed_out: false,
per_disk: vec![
DiskCapacityUpdate {
disk: CapacityScopeDisk {
+4
View File
@@ -28,6 +28,10 @@ pub(crate) struct CapacityScanResult {
pub is_estimated: bool,
pub scan_duration: Duration,
pub had_partial_errors: bool,
pub timed_out: bool,
pub metadata_incomplete: bool,
#[cfg(test)]
pub metadata_reads: usize,
}
impl CapacityScanResult {