mirror of
https://github.com/rustfs/rustfs.git
synced 2026-09-05 11:45:39 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 43a7f94f14 |
@@ -230,7 +230,7 @@ fn scanner_abandoned_child_list_options() -> ListPathRawOptions {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub fn data_usage_update_dir_cycles() -> u32 {
|
pub fn data_usage_update_dir_cycles() -> u32 {
|
||||||
rustfs_utils::get_env_u32(ENV_DATA_USAGE_UPDATE_DIR_CYCLES, DATA_USAGE_UPDATE_DIR_CYCLES)
|
rustfs_utils::get_env_u32(ENV_DATA_USAGE_UPDATE_DIR_CYCLES, DATA_USAGE_UPDATE_DIR_CYCLES).max(1)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn heal_object_select_prob() -> u32 {
|
pub fn heal_object_select_prob() -> u32 {
|
||||||
@@ -806,6 +806,7 @@ impl FolderScanner {
|
|||||||
fn prune_failed_objects_cache(&mut self) {
|
fn prune_failed_objects_cache(&mut self) {
|
||||||
let ttl = self.failed_object_ttl_secs;
|
let ttl = self.failed_object_ttl_secs;
|
||||||
if ttl == 0 {
|
if ttl == 0 {
|
||||||
|
self.new_cache.info.failed_objects.clear();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -963,6 +964,27 @@ impl FolderScanner {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn preserve_failed_child(
|
||||||
|
&mut self,
|
||||||
|
parent: &Option<DataUsageHash>,
|
||||||
|
child_hash: &DataUsageHash,
|
||||||
|
parent_entry: &mut DataUsageEntry,
|
||||||
|
child_entry: &DataUsageEntry,
|
||||||
|
) {
|
||||||
|
// A failed walk proves neither deletion nor a complete replacement.
|
||||||
|
// Keep the previous subtree and mark this snapshot incomplete even
|
||||||
|
// when the failed-object retry cache is disabled or at capacity.
|
||||||
|
parent_entry.failed_objects = parent_entry.failed_objects.saturating_add(1);
|
||||||
|
if self.old_cache.cache.contains_key(&child_hash.key()) {
|
||||||
|
self.new_cache.delete_recursive(child_hash);
|
||||||
|
self.new_cache.copy_with_children(&self.old_cache, child_hash, parent);
|
||||||
|
parent_entry.add_child(child_hash);
|
||||||
|
} else {
|
||||||
|
self.preserve_partial_child_progress(parent, child_hash, parent_entry, child_entry)
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fn alert_excessive_folders(&self, folder: &str, total_folders: usize) {
|
fn alert_excessive_folders(&self, folder: &str, total_folders: usize) {
|
||||||
let threshold = scanner_excess_folders_threshold();
|
let threshold = scanner_excess_folders_threshold();
|
||||||
if u64::try_from(total_folders).unwrap_or(u64::MAX) <= threshold {
|
if u64::try_from(total_folders).unwrap_or(u64::MAX) <= threshold {
|
||||||
@@ -1177,8 +1199,6 @@ impl FolderScanner {
|
|||||||
return Err(ScannerError::Other("Operation cancelled".to_string()));
|
return Err(ScannerError::Other("Operation cancelled".to_string()));
|
||||||
}
|
}
|
||||||
|
|
||||||
self.prune_failed_objects_cache();
|
|
||||||
|
|
||||||
let mut abandoned_children: DataUsageHashMap = HashSet::new();
|
let mut abandoned_children: DataUsageHashMap = HashSet::new();
|
||||||
if !into.compacted {
|
if !into.compacted {
|
||||||
abandoned_children = self.old_cache.find_children_copy(this_hash.clone());
|
abandoned_children = self.old_cache.find_children_copy(this_hash.clone());
|
||||||
@@ -1221,7 +1241,9 @@ impl FolderScanner {
|
|||||||
};
|
};
|
||||||
let active_object_lock = self.old_cache.info.object_lock.clone();
|
let active_object_lock = self.old_cache.info.object_lock.clone();
|
||||||
|
|
||||||
self.sleeper.sleep_folder().await;
|
ctx.run_until_cancelled(self.sleeper.sleep_folder())
|
||||||
|
.await
|
||||||
|
.ok_or_else(|| ScannerError::Other("Operation cancelled".to_string()))?;
|
||||||
|
|
||||||
let mut existing_folders: Vec<CachedFolder> = Vec::new();
|
let mut existing_folders: Vec<CachedFolder> = Vec::new();
|
||||||
let mut new_folders: Vec<CachedFolder> = Vec::new();
|
let mut new_folders: Vec<CachedFolder> = Vec::new();
|
||||||
@@ -1448,7 +1470,7 @@ impl FolderScanner {
|
|||||||
|
|
||||||
let heal_enabled = this_hash.mod_alt(
|
let heal_enabled = this_hash.mod_alt(
|
||||||
self.old_cache.info.next_cycle as u32 / folder.object_heal_prob_div,
|
self.old_cache.info.next_cycle as u32 / folder.object_heal_prob_div,
|
||||||
self.heal_object_select / folder.object_heal_prob_div,
|
(self.heal_object_select / folder.object_heal_prob_div).max(1),
|
||||||
) && self.should_heal().await;
|
) && self.should_heal().await;
|
||||||
|
|
||||||
let mut item = ScannerItem {
|
let mut item = ScannerItem {
|
||||||
@@ -1465,12 +1487,10 @@ impl FolderScanner {
|
|||||||
file_type: entry_type,
|
file_type: entry_type,
|
||||||
};
|
};
|
||||||
|
|
||||||
// If this path is already known as failed, just skip it.
|
// Count unresolved objects in each snapshot without extending
|
||||||
// We intentionally do NOT call `record_failed` or bump `failed_objects` here,
|
// the retry TTL or emitting another failure event.
|
||||||
// because the failure was recorded when the original error occurred
|
|
||||||
// (e.g. in the get_size error branch below). This branch only accounts
|
|
||||||
// for subsequent skips of already-failed paths.
|
|
||||||
if self.should_skip_failed(&item.path) {
|
if self.should_skip_failed(&item.path) {
|
||||||
|
into.failed_objects = into.failed_objects.saturating_add(1);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1485,7 +1505,7 @@ impl FolderScanner {
|
|||||||
|
|
||||||
if failure_action != GetSizeFailureAction::Skip {
|
if failure_action != GetSizeFailureAction::Skip {
|
||||||
// Track failed objects to prevent infinite retry loops
|
// Track failed objects to prevent infinite retry loops
|
||||||
into.failed_objects += 1;
|
into.failed_objects = into.failed_objects.saturating_add(1);
|
||||||
self.record_failed(&item.path);
|
self.record_failed(&item.path);
|
||||||
|
|
||||||
if should_log_failed_object(into.failed_objects) {
|
if should_log_failed_object(into.failed_objects) {
|
||||||
@@ -1564,12 +1584,15 @@ impl FolderScanner {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
timer.sleep().await;
|
ctx.run_until_cancelled(timer.sleep())
|
||||||
|
.await
|
||||||
|
.ok_or_else(|| ScannerError::Other("Operation cancelled".to_string()))?;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
found_object_metadata = true;
|
found_object_metadata = true;
|
||||||
|
self.new_cache.info.failed_objects.remove(&item.path);
|
||||||
|
|
||||||
item.transform_meta_dir();
|
item.transform_meta_dir();
|
||||||
|
|
||||||
@@ -1581,7 +1604,9 @@ impl FolderScanner {
|
|||||||
object_count += 1;
|
object_count += 1;
|
||||||
self.budget.record_object_scanned();
|
self.budget.record_object_scanned();
|
||||||
|
|
||||||
timer.sleep().await;
|
ctx.run_until_cancelled(timer.sleep())
|
||||||
|
.await
|
||||||
|
.ok_or_else(|| ScannerError::Other("Operation cancelled".to_string()))?;
|
||||||
|
|
||||||
if ctx.is_cancelled() {
|
if ctx.is_cancelled() {
|
||||||
return Err(ScannerError::Other("Operation cancelled".to_string()));
|
return Err(ScannerError::Other("Operation cancelled".to_string()));
|
||||||
@@ -1622,9 +1647,9 @@ impl FolderScanner {
|
|||||||
if self.is_erasure_mode && found_erasure_data_directory && !found_object_metadata {
|
if self.is_erasure_mode && found_erasure_data_directory && !found_object_metadata {
|
||||||
found_object_metadata = true;
|
found_object_metadata = true;
|
||||||
let metadata_path = path_join_buf(&[&dir_path, STORAGE_FORMAT_FILE]);
|
let metadata_path = path_join_buf(&[&dir_path, STORAGE_FORMAT_FILE]);
|
||||||
|
into.failed_objects = into.failed_objects.saturating_add(1);
|
||||||
|
|
||||||
if !self.should_skip_failed(&metadata_path) {
|
if !self.should_skip_failed(&metadata_path) {
|
||||||
into.failed_objects = into.failed_objects.saturating_add(1);
|
|
||||||
self.record_failed(&metadata_path);
|
self.record_failed(&metadata_path);
|
||||||
|
|
||||||
let failed_cache_entries = self.new_cache.info.failed_objects.len();
|
let failed_cache_entries = self.new_cache.info.failed_objects.len();
|
||||||
@@ -1835,6 +1860,7 @@ impl FolderScanner {
|
|||||||
error = %e,
|
error = %e,
|
||||||
"Scanner child folder scan failed"
|
"Scanner child folder scan failed"
|
||||||
);
|
);
|
||||||
|
self.preserve_failed_child(&folder_item.parent, &h, into, &dst).await;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
tokio::task::yield_now().await;
|
tokio::task::yield_now().await;
|
||||||
@@ -2230,6 +2256,7 @@ impl FolderScanner {
|
|||||||
error = %e,
|
error = %e,
|
||||||
"Scanner heal child folder scan failed"
|
"Scanner heal child folder scan failed"
|
||||||
);
|
);
|
||||||
|
self.preserve_failed_child(&folder_item.parent, &h, into, &dst).await;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
tokio::task::yield_now().await;
|
tokio::task::yield_now().await;
|
||||||
@@ -2396,6 +2423,9 @@ pub async fn scan_data_folder(
|
|||||||
};
|
};
|
||||||
|
|
||||||
let now = FolderScanner::now_secs();
|
let now = FolderScanner::now_secs();
|
||||||
|
// Prune once per bucket walk, not once per directory. Per-path TTL checks
|
||||||
|
// still allow retries during long scans, and insertions enforce the cap.
|
||||||
|
scanner.prune_failed_objects_cache();
|
||||||
prune_size_reconciliation(&mut scanner.new_cache.info, now);
|
prune_size_reconciliation(&mut scanner.new_cache.info, now);
|
||||||
prune_size_reconciliation(&mut scanner.update_cache.info, now);
|
prune_size_reconciliation(&mut scanner.update_cache.info, now);
|
||||||
|
|
||||||
@@ -2422,7 +2452,9 @@ pub async fn scan_data_folder(
|
|||||||
new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN);
|
new_cache.force_compact(DATA_SCANNER_COMPACT_AT_CHILDREN);
|
||||||
new_cache.info.last_update = Some(SystemTime::now());
|
new_cache.info.last_update = Some(SystemTime::now());
|
||||||
new_cache.info.next_cycle = cache.info.next_cycle;
|
new_cache.info.next_cycle = cache.info.next_cycle;
|
||||||
let unresolved_objects = root.failed_objects > 0
|
let unresolved_objects = new_cache
|
||||||
|
.size_recursive(&cache.info.name)
|
||||||
|
.is_none_or(|root| root.failed_objects > 0)
|
||||||
|| !new_cache.info.failed_objects.is_empty()
|
|| !new_cache.info.failed_objects.is_empty()
|
||||||
|| !new_cache.info.size_reconciliation.is_empty();
|
|| !new_cache.info.size_reconciliation.is_empty();
|
||||||
new_cache.info.snapshot_complete = !unresolved_objects;
|
new_cache.info.snapshot_complete = !unresolved_objects;
|
||||||
|
|||||||
@@ -1366,6 +1366,104 @@ mod tests {
|
|||||||
assert_eq!(item.object_path(), "object");
|
assert_eq!(item.object_path(), "object");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial_test::serial]
|
||||||
|
async fn scanner_blocked_expiry_preserves_usage_replication_and_integrity_work() {
|
||||||
|
use s3s::dto::{LifecycleExpiration, LifecycleRule};
|
||||||
|
|
||||||
|
let lifecycle = Arc::new(BucketLifecycleConfiguration {
|
||||||
|
rules: vec![LifecycleRule {
|
||||||
|
status: ExpirationStatus::from_static(ExpirationStatus::ENABLED),
|
||||||
|
expiration: Some(LifecycleExpiration {
|
||||||
|
days: Some(1),
|
||||||
|
..Default::default()
|
||||||
|
}),
|
||||||
|
abort_incomplete_multipart_upload: None,
|
||||||
|
del_marker_expiration: None,
|
||||||
|
filter: None,
|
||||||
|
id: None,
|
||||||
|
noncurrent_version_expiration: None,
|
||||||
|
noncurrent_version_transitions: None,
|
||||||
|
prefix: None,
|
||||||
|
transitions: None,
|
||||||
|
}],
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
let attempts = |report: &rustfs_scanner_metrics::metrics::ScannerMetricsReport, source: ScannerWorkSource| {
|
||||||
|
report
|
||||||
|
.source_work
|
||||||
|
.iter()
|
||||||
|
.filter(|work| work.source == source.as_str())
|
||||||
|
.map(|work| work.queued + work.skipped + work.missed)
|
||||||
|
.sum::<u64>()
|
||||||
|
};
|
||||||
|
for with_lifecycle in [false, true] {
|
||||||
|
for scan_mode in [HealScanMode::Normal, HealScanMode::Deep] {
|
||||||
|
for guard in ["pending", "failed", "legal_hold"] {
|
||||||
|
let mut metadata = HashMap::new();
|
||||||
|
let replication_status = match guard {
|
||||||
|
"pending" => ReplicationStatusType::Pending,
|
||||||
|
"failed" => ReplicationStatusType::Failed,
|
||||||
|
_ => {
|
||||||
|
metadata.insert("x-amz-object-lock-legal-hold".to_string(), "ON".to_string());
|
||||||
|
ReplicationStatusType::Completed
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let object = ObjectInfo {
|
||||||
|
bucket: "bucket".to_string(),
|
||||||
|
name: "object".to_string(),
|
||||||
|
version_id: Some(uuid::Uuid::new_v4()),
|
||||||
|
num_versions: 1,
|
||||||
|
is_latest: true,
|
||||||
|
mod_time: Some(OffsetDateTime::now_utc() - time::Duration::days(90)),
|
||||||
|
size: 4096,
|
||||||
|
actual_size: 4096,
|
||||||
|
replication_status,
|
||||||
|
user_defined: Arc::new(metadata),
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
let events = Evaluator::new(lifecycle.clone())
|
||||||
|
.eval(&[crate::ecstore_object_opts_from_object_info(&object)])
|
||||||
|
.await
|
||||||
|
.expect("evaluate expiry guard");
|
||||||
|
assert_eq!(events[0].action, IlmAction::NoneAction, "expiry must be blocked by {guard}");
|
||||||
|
|
||||||
|
let mut item = scanner_item_with_prefix("");
|
||||||
|
item.object_name = "object".to_string();
|
||||||
|
item.lifecycle = with_lifecycle.then(|| lifecycle.clone());
|
||||||
|
item.replication = Some(Arc::new(ReplicationConfig::new(None, None)));
|
||||||
|
item.heal_enabled = true;
|
||||||
|
item.heal_bitrot = scan_mode == HealScanMode::Deep;
|
||||||
|
let before = global_metrics().report().await;
|
||||||
|
let mut summary = SizeSummary::default();
|
||||||
|
item.apply_actions(vec![object], None, VersioningConfiguration::default(), &[], &mut summary)
|
||||||
|
.await;
|
||||||
|
let after = global_metrics().report().await;
|
||||||
|
assert_eq!(summary.total_size, 4096, "blocked expiry must retain bytes for {guard}");
|
||||||
|
assert_eq!(summary.versions, 1);
|
||||||
|
assert_eq!(summary.delete_markers, 0);
|
||||||
|
assert!(summary.size_reconciliation.is_empty());
|
||||||
|
assert_eq!(
|
||||||
|
attempts(&after, scanner_heal_source(scan_mode)) - attempts(&before, scanner_heal_source(scan_mode)),
|
||||||
|
1,
|
||||||
|
"integrity work must continue with lifecycle={with_lifecycle}, guard={guard}"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
attempts(&after, ScannerWorkSource::BucketReplication)
|
||||||
|
- attempts(&before, ScannerWorkSource::BucketReplication),
|
||||||
|
1,
|
||||||
|
"replication inspection must continue with lifecycle={with_lifecycle}, guard={guard}"
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
attempts(&after, ScannerWorkSource::Lifecycle) - attempts(&before, ScannerWorkSource::Lifecycle),
|
||||||
|
0,
|
||||||
|
"blocked expiry must not enqueue destructive lifecycle work"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn unknown_tier_never_triggers_transition() {
|
fn unknown_tier_never_triggers_transition() {
|
||||||
let object = ObjectInfo {
|
let object = ObjectInfo {
|
||||||
|
|||||||
@@ -1883,6 +1883,323 @@ async fn test_scan_folder_skips_unreadable_child_directory() {
|
|||||||
assert!(result.is_ok(), "expected unreadable child directory to be skipped");
|
assert!(result.is_ok(), "expected unreadable child directory to be skipped");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn scanner_failed_child_retains_usage_and_scans_healthy_sibling() {
|
||||||
|
for with_prior in [false, true] {
|
||||||
|
let (mut scanner, temp_dir) = build_test_scanner().await;
|
||||||
|
let _guard = TestGuard::new(0, 0, &mut scanner, temp_dir.clone());
|
||||||
|
let bad_dir = temp_dir.join("bucket/bad");
|
||||||
|
tokio::fs::create_dir_all(&bad_dir).await.expect("create failing directory");
|
||||||
|
write_test_object_metadata_bytes(
|
||||||
|
&temp_dir,
|
||||||
|
"bucket",
|
||||||
|
"good",
|
||||||
|
&metadata_for_object_version("bucket", "good", Some(Uuid::new_v4())),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
scanner.old_cache.info.name = "bucket".to_string();
|
||||||
|
scanner.new_cache.info.name = "bucket".to_string();
|
||||||
|
scanner.update_cache.info.name = "bucket".to_string();
|
||||||
|
let root_hash = hash_path("bucket");
|
||||||
|
let bad_hash = hash_path("bucket/bad");
|
||||||
|
let mut prior = DataUsageEntry {
|
||||||
|
size: 4096,
|
||||||
|
objects: 2,
|
||||||
|
versions: 3,
|
||||||
|
delete_markers: 1,
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
prior.replication_stats = Some(rustfs_data_usage::ReplicationAllStats {
|
||||||
|
replica_size: 4096,
|
||||||
|
replica_count: 2,
|
||||||
|
..Default::default()
|
||||||
|
});
|
||||||
|
prior.add_tier_sizes(&HashMap::from([(
|
||||||
|
"WARM".to_string(),
|
||||||
|
TierStats {
|
||||||
|
total_size: 4096,
|
||||||
|
num_versions: 3,
|
||||||
|
num_objects: 2,
|
||||||
|
},
|
||||||
|
)]));
|
||||||
|
scanner
|
||||||
|
.old_cache
|
||||||
|
.replace_hashed(&root_hash, &None, &DataUsageEntry::default());
|
||||||
|
if with_prior {
|
||||||
|
scanner.old_cache.replace_hashed(&bad_hash, &Some(root_hash.clone()), &prior);
|
||||||
|
} else {
|
||||||
|
prior = DataUsageEntry::default();
|
||||||
|
}
|
||||||
|
scanner.update_current_path = Arc::new(move |path| {
|
||||||
|
if path == "bucket/bad" {
|
||||||
|
// Replace the directory after enumeration but before descent. This
|
||||||
|
// injects a real read_dir error even when tests run as root.
|
||||||
|
std::fs::remove_dir(&bad_dir).expect("remove enumerated directory");
|
||||||
|
std::fs::write(&bad_dir, b"not a directory").expect("replace enumerated directory");
|
||||||
|
}
|
||||||
|
Box::pin(async {})
|
||||||
|
});
|
||||||
|
let mut root = DataUsageEntry::default();
|
||||||
|
scanner
|
||||||
|
.scan_folder(
|
||||||
|
CancellationToken::new(),
|
||||||
|
CachedFolder {
|
||||||
|
name: "bucket".to_string(),
|
||||||
|
parent: None,
|
||||||
|
object_heal_prob_div: 1,
|
||||||
|
},
|
||||||
|
&mut root,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("one failed directory must not stop healthy siblings");
|
||||||
|
let total = scanner.new_cache.size_recursive(&root_hash.key()).expect("root usage");
|
||||||
|
assert_eq!(total.size, prior.size + 1, "unreadable child must retain its previous bytes");
|
||||||
|
assert_eq!(total.objects, prior.objects + 1, "healthy sibling must still be counted");
|
||||||
|
assert_eq!(total.versions, prior.versions + 1);
|
||||||
|
assert_eq!(total.delete_markers, prior.delete_markers);
|
||||||
|
assert_eq!(total.failed_objects, 1, "walk error must keep the snapshot incomplete");
|
||||||
|
assert_eq!(
|
||||||
|
serde_json::to_value(&total.replication_stats).expect("serialize replication usage"),
|
||||||
|
serde_json::to_value(&prior.replication_stats).expect("serialize prior replication usage")
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
serde_json::to_value(&total.all_tier_stats).expect("serialize tier usage"),
|
||||||
|
serde_json::to_value(&prior.all_tier_stats).expect("serialize prior tier usage")
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn scanner_nested_metadata_failure_without_retry_cache_is_partial_then_recovers() {
|
||||||
|
let (scanner, temp_dir) = build_test_scanner().await;
|
||||||
|
let _guard = TestGuard {
|
||||||
|
temp_dir: Some(temp_dir.clone()),
|
||||||
|
};
|
||||||
|
write_test_object_metadata_bytes(&temp_dir, "bucket", "prefix/bad", b"").await;
|
||||||
|
write_test_object_metadata_bytes(
|
||||||
|
&temp_dir,
|
||||||
|
"bucket",
|
||||||
|
"prefix/good",
|
||||||
|
&metadata_for_object_version("bucket", "prefix/good", Some(Uuid::new_v4())),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
temp_env::async_with_vars([(ENV_FAILED_OBJECT_TTL_SECS, Some("0"))], async {
|
||||||
|
for inherited_failure in [false, true] {
|
||||||
|
write_test_object_metadata_bytes(&temp_dir, "bucket", "prefix/bad", b"").await;
|
||||||
|
let mut cache = DataUsageCache {
|
||||||
|
info: DataUsageCacheInfo {
|
||||||
|
name: "bucket".to_string(),
|
||||||
|
next_cycle: u64::from(
|
||||||
|
(0..DATA_USAGE_UPDATE_DIR_CYCLES)
|
||||||
|
.find(|cycle| !hash_path("bucket/prefix").mod_(*cycle, DATA_USAGE_UPDATE_DIR_CYCLES))
|
||||||
|
.expect("cycle outside the prefix compaction sample"),
|
||||||
|
),
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
..Default::default()
|
||||||
|
};
|
||||||
|
if inherited_failure {
|
||||||
|
cache
|
||||||
|
.info
|
||||||
|
.failed_objects
|
||||||
|
.insert("removed-object/xl.meta".to_string(), FolderScanner::now_secs());
|
||||||
|
}
|
||||||
|
let budget = ScannerCycleBudget::new(&CancellationToken::new(), Default::default());
|
||||||
|
let result = scan_data_folder(
|
||||||
|
budget.token(),
|
||||||
|
budget,
|
||||||
|
vec![scanner.local_disk.clone()],
|
||||||
|
scanner.local_disk.clone(),
|
||||||
|
cache,
|
||||||
|
None,
|
||||||
|
HealScanMode::Normal,
|
||||||
|
DynamicSleeper::new(rustfs_config::ScannerSpeed::Fastest),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
let mut partial = match result {
|
||||||
|
Err(ScannerError::PartialCache(cache)) => *cache,
|
||||||
|
other => panic!("nested failure must never publish a complete snapshot: {other:?}"),
|
||||||
|
};
|
||||||
|
assert!(!partial.info.snapshot_complete);
|
||||||
|
assert!(partial.info.failed_objects.is_empty(), "TTL zero disables only the retry cache");
|
||||||
|
let total = partial.size_recursive("bucket").expect("partial root");
|
||||||
|
assert_eq!(total.objects, 1);
|
||||||
|
assert_eq!(total.failed_objects, 1);
|
||||||
|
|
||||||
|
// Reusing the partial compacted subtree must remain partial even
|
||||||
|
// without a retry ledger. Recovery happens on its next selected cycle.
|
||||||
|
let budget = ScannerCycleBudget::new(&CancellationToken::new(), Default::default());
|
||||||
|
let reused = scan_data_folder(
|
||||||
|
budget.token(),
|
||||||
|
budget,
|
||||||
|
vec![scanner.local_disk.clone()],
|
||||||
|
scanner.local_disk.clone(),
|
||||||
|
partial.clone(),
|
||||||
|
None,
|
||||||
|
HealScanMode::Normal,
|
||||||
|
DynamicSleeper::new(rustfs_config::ScannerSpeed::Fastest),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
assert!(matches!(reused, Err(ScannerError::PartialCache(_))));
|
||||||
|
partial.info.next_cycle = u64::from(
|
||||||
|
(0..DATA_USAGE_UPDATE_DIR_CYCLES)
|
||||||
|
.find(|cycle| hash_path("bucket/prefix").mod_(*cycle, DATA_USAGE_UPDATE_DIR_CYCLES))
|
||||||
|
.expect("next selected directory cycle"),
|
||||||
|
);
|
||||||
|
write_test_object_metadata_bytes(
|
||||||
|
&temp_dir,
|
||||||
|
"bucket",
|
||||||
|
"prefix/bad",
|
||||||
|
&metadata_for_object_version("bucket", "prefix/bad", Some(Uuid::new_v4())),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
let budget = ScannerCycleBudget::new(&CancellationToken::new(), Default::default());
|
||||||
|
let recovered = scan_data_folder(
|
||||||
|
budget.token(),
|
||||||
|
budget,
|
||||||
|
vec![scanner.local_disk.clone()],
|
||||||
|
scanner.local_disk.clone(),
|
||||||
|
partial,
|
||||||
|
None,
|
||||||
|
HealScanMode::Normal,
|
||||||
|
DynamicSleeper::new(rustfs_config::ScannerSpeed::Fastest),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("repaired subtree must converge on its next selected cycle");
|
||||||
|
assert!(recovered.info.snapshot_complete);
|
||||||
|
let total = recovered.size_recursive("bucket").expect("recovered root");
|
||||||
|
assert_eq!(total.objects, 2);
|
||||||
|
assert_eq!(total.size, 2);
|
||||||
|
assert_eq!(total.versions, 2);
|
||||||
|
assert_eq!(total.failed_objects, 0);
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn scanner_compacted_directory_keeps_aggressive_heal_and_bitrot_sampling() {
|
||||||
|
for scan_mode in [HealScanMode::Normal, HealScanMode::Deep] {
|
||||||
|
for select_prob in [0, 1, 8, 16] {
|
||||||
|
let (mut scanner, temp_dir) = build_test_scanner().await;
|
||||||
|
let _guard = TestGuard::new(0, 0, &mut scanner, temp_dir.clone());
|
||||||
|
write_test_object_metadata_bytes(
|
||||||
|
&temp_dir,
|
||||||
|
"bucket",
|
||||||
|
"object",
|
||||||
|
&metadata_for_object_version("bucket", "object", Some(Uuid::new_v4())),
|
||||||
|
)
|
||||||
|
.await;
|
||||||
|
scanner.old_cache.info.name = "bucket".to_string();
|
||||||
|
scanner.new_cache.info.name = "bucket".to_string();
|
||||||
|
scanner.update_cache.info.name = "bucket".to_string();
|
||||||
|
scanner.is_erasure_mode = true;
|
||||||
|
scanner.heal_object_select = select_prob;
|
||||||
|
scanner.scan_mode = scan_mode;
|
||||||
|
let root_hash = hash_path("bucket");
|
||||||
|
let object_hash = hash_path("bucket/object");
|
||||||
|
scanner.old_cache.info.next_cycle = u64::from(
|
||||||
|
(0..DATA_USAGE_UPDATE_DIR_CYCLES)
|
||||||
|
.find(|cycle| object_hash.mod_(*cycle, DATA_USAGE_UPDATE_DIR_CYCLES))
|
||||||
|
.expect("selected directory cycle"),
|
||||||
|
);
|
||||||
|
scanner
|
||||||
|
.old_cache
|
||||||
|
.replace_hashed(&root_hash, &None, &DataUsageEntry::default());
|
||||||
|
scanner.old_cache.replace_hashed(
|
||||||
|
&object_hash,
|
||||||
|
&Some(root_hash),
|
||||||
|
&DataUsageEntry {
|
||||||
|
compacted: true,
|
||||||
|
objects: 1,
|
||||||
|
versions: 1,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
);
|
||||||
|
let attempts = |report: rustfs_scanner_metrics::metrics::ScannerMetricsReport| {
|
||||||
|
report
|
||||||
|
.source_work
|
||||||
|
.iter()
|
||||||
|
.filter(|work| work.source == scanner_heal_source(scan_mode).as_str())
|
||||||
|
.map(|work| work.queued + work.skipped + work.missed)
|
||||||
|
.sum::<u64>()
|
||||||
|
};
|
||||||
|
temp_env::async_with_vars([(ENV_SCANNER_DEEP_VERIFY_COOLDOWN_SECS, Some("0"))], async {
|
||||||
|
let before = attempts(global_metrics().report().await);
|
||||||
|
let mut root = DataUsageEntry::default();
|
||||||
|
scanner
|
||||||
|
.scan_folder(
|
||||||
|
CancellationToken::new(),
|
||||||
|
CachedFolder {
|
||||||
|
name: "bucket".to_string(),
|
||||||
|
parent: None,
|
||||||
|
object_heal_prob_div: 1,
|
||||||
|
},
|
||||||
|
&mut root,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("scan selected compacted object");
|
||||||
|
assert_eq!(
|
||||||
|
attempts(global_metrics().report().await) - before,
|
||||||
|
u64::from(select_prob != 0),
|
||||||
|
"selected compacted object must reach {scan_mode:?} admission with divisor {select_prob}"
|
||||||
|
);
|
||||||
|
let total = scanner.new_cache.size_recursive("bucket").expect("usage root");
|
||||||
|
assert_eq!(total.objects, 1);
|
||||||
|
assert_eq!(total.versions, 1);
|
||||||
|
})
|
||||||
|
.await;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test(start_paused = true)]
|
||||||
|
#[serial]
|
||||||
|
async fn scanner_cancellation_interrupts_folder_throttle() {
|
||||||
|
use futures::{FutureExt, poll};
|
||||||
|
let (mut scanner, temp_dir) = build_test_scanner().await;
|
||||||
|
let _guard = TestGuard::new(0, 0, &mut scanner, temp_dir);
|
||||||
|
scanner.sleeper = DynamicSleeper::new(rustfs_config::ScannerSpeed::Slowest);
|
||||||
|
let previous_idle = crate::sleeper::SCANNER_IDLE_MODE.swap(true, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
let ctx = CancellationToken::new();
|
||||||
|
let mut root = DataUsageEntry::default();
|
||||||
|
let mut scan = scanner
|
||||||
|
.scan_folder(
|
||||||
|
ctx.clone(),
|
||||||
|
CachedFolder {
|
||||||
|
name: "bucket".to_string(),
|
||||||
|
parent: None,
|
||||||
|
object_heal_prob_div: 1,
|
||||||
|
},
|
||||||
|
&mut root,
|
||||||
|
)
|
||||||
|
.boxed();
|
||||||
|
assert!(poll!(scan.as_mut()).is_pending(), "scan should be waiting in its folder throttle");
|
||||||
|
ctx.cancel();
|
||||||
|
let outcome = scan.now_or_never();
|
||||||
|
crate::sleeper::SCANNER_IDLE_MODE.store(previous_idle, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
assert!(
|
||||||
|
matches!(outcome, Some(Err(_))),
|
||||||
|
"cancellation must finish without advancing the sleep clock"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
#[serial]
|
||||||
|
fn scanner_zero_directory_cycle_keeps_rescanning_enabled() {
|
||||||
|
temp_env::with_var(ENV_DATA_USAGE_UPDATE_DIR_CYCLES, Some("0"), || {
|
||||||
|
for cycle in 0..32 {
|
||||||
|
assert!(
|
||||||
|
hash_path("bucket/object").mod_(cycle, data_usage_update_dir_cycles()),
|
||||||
|
"zero must not leave compacted usage stale forever"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
#[serial]
|
#[serial]
|
||||||
async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
|
async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
|
||||||
@@ -2153,7 +2470,7 @@ async fn test_scan_folder_corrupt_xl_meta_stops_erasure_data_dir_descent() {
|
|||||||
.await
|
.await
|
||||||
.expect("cached metadata failure must still stop erasure data directory descent");
|
.expect("cached metadata failure must still stop erasure data directory descent");
|
||||||
|
|
||||||
assert_eq!(retry_into.failed_objects, 0, "cached failure should not be counted twice");
|
assert_eq!(retry_into.failed_objects, 1, "cached failure must remain visible in each snapshot");
|
||||||
assert!(!retry_budget.budget_elapsed());
|
assert!(!retry_budget.budget_elapsed());
|
||||||
assert_eq!(retry_budget.reason(), None);
|
assert_eq!(retry_budget.reason(), None);
|
||||||
|
|
||||||
@@ -2278,7 +2595,7 @@ async fn test_scan_folder_missing_xl_meta_stops_erasure_data_dir_descent() {
|
|||||||
.await
|
.await
|
||||||
.expect("cached missing metadata must still stop erasure data directory descent");
|
.expect("cached missing metadata must still stop erasure data directory descent");
|
||||||
|
|
||||||
assert_eq!(retry_into.failed_objects, 0, "cached failure should not be counted twice");
|
assert_eq!(retry_into.failed_objects, 1, "cached failure must remain visible in each snapshot");
|
||||||
assert!(!retry_budget.budget_elapsed());
|
assert!(!retry_budget.budget_elapsed());
|
||||||
assert_eq!(retry_budget.reason(), None);
|
assert_eq!(retry_budget.reason(), None);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -34,6 +34,22 @@ The `scanner` and `heal` subsystems are served by `GetConfigKVHandler` (`rustfs/
|
|||||||
|
|
||||||
## Test Matrix
|
## Test Matrix
|
||||||
|
|
||||||
|
### Deterministic regression checks
|
||||||
|
|
||||||
|
Run the scanner regressions before collecting host-pressure measurements:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cargo nextest run -p rustfs-scanner --lib
|
||||||
|
```
|
||||||
|
|
||||||
|
Most tests in `crates/scanner/tests/lifecycle_integration_test.rs` are ignored in the default lane because they require serial execution. Run the scanner portion of the `ILM Integration (serial)` selection in `.github/workflows/ci.yml` with `-j1 --run-ignored all` as well; preserve its documented exclusions for known noncurrent transition/expiry failures.
|
||||||
|
|
||||||
|
The folder regressions exercise real directory enumeration and metadata decoding. `scanner_failed_child_retains_usage_and_scans_healthy_sibling` replaces an enumerated directory before descent, so its I/O failure is reproducible without depending on Unix permission enforcement. `scanner_nested_metadata_failure_without_retry_cache_is_partial_then_recovers` checks fresh and inherited failure state with retry caching disabled, reuse of a partial compacted subtree, and recovery on the next selected directory cycle. Neither a failed subtree nor an expired retry ledger proves zero usage.
|
||||||
|
|
||||||
|
`scanner_compacted_directory_keeps_aggressive_heal_and_bitrot_sampling` covers disabled, sub-interval, and exact-interval heal divisors in normal and deep modes. `scanner_cancellation_interrupts_folder_throttle` uses a paused clock to require immediate cooperative cancellation. `scanner_blocked_expiry_preserves_usage_replication_and_integrity_work` covers lifecycle enabled/disabled with pending replication, failed replication, and Legal Hold; retained bytes and integrity/replication inspection must survive blocked expiry.
|
||||||
|
|
||||||
|
These checks complement the sampling and cancellation design in [MinIO's scanner implementation](https://github.com/minio/minio/blob/master/cmd/data-scanner.go), especially `scanDataFolder`, `folderScanner.scanFolder`, and `dynamicSleeper.Sleep`. Scanner admission counters prove that work reaches the admission boundary; they do not prove remote replication delivery or a completed shard repair. The deployment matrix below remains necessary for those claims and for measured CPU, memory, IOPS, and foreground-latency comparisons.
|
||||||
|
|
||||||
Collect at least two runs on the same RustFS commit and the same workload. Keep hardware, commit, object count, object size, bucket count, scanner-enabled state, and foreground workload constant between runs.
|
Collect at least two runs on the same RustFS commit and the same workload. Keep hardware, commit, object count, object size, bucket count, scanner-enabled state, and foreground workload constant between runs.
|
||||||
|
|
||||||
| Run | Purpose | Example scanner settings |
|
| Run | Purpose | Example scanner settings |
|
||||||
|
|||||||
@@ -70,9 +70,9 @@ These have no persistent key and are read from the environment only.
|
|||||||
| `RUSTFS_SCANNER_ENABLED` (deprecated alias `RUSTFS_ENABLE_SCANNER`) | `true` (`scanner_enabled_from_env`, `rustfs/src/module_switches.rs`) | Starts the data scanner at all. The heal manager is initialized whenever heal or scanner is enabled, because scanner-produced heal candidates need a consumer. |
|
| `RUSTFS_SCANNER_ENABLED` (deprecated alias `RUSTFS_ENABLE_SCANNER`) | `true` (`scanner_enabled_from_env`, `rustfs/src/module_switches.rs`) | Starts the data scanner at all. The heal manager is initialized whenever heal or scanner is enabled, because scanner-produced heal candidates need a consumer. |
|
||||||
| `RUSTFS_SCANNER_ALERT_COOLDOWN_SECS` | `86400` (`DEFAULT_SCANNER_ALERT_COOLDOWN_SECS`, `scanner_folder.rs`) | Per-(kind, bucket, object) cooldown between S3 excess-alert events; `0` emits every cycle. See [Scanner Excess Alerts](scanner-excess-alerts.md). |
|
| `RUSTFS_SCANNER_ALERT_COOLDOWN_SECS` | `86400` (`DEFAULT_SCANNER_ALERT_COOLDOWN_SECS`, `scanner_folder.rs`) | Per-(kind, bucket, object) cooldown between S3 excess-alert events; `0` emits every cycle. See [Scanner Excess Alerts](scanner-excess-alerts.md). |
|
||||||
| `RUSTFS_SCANNER_DEEP_VERIFY_COOLDOWN_SECS` | `60` (`DEFAULT_SCANNER_DEEP_VERIFY_COOLDOWN_SECS`, `scanner_folder.rs`) | Objects modified within this window are skipped by deep (bitrot) verification in the current cycle. |
|
| `RUSTFS_SCANNER_DEEP_VERIFY_COOLDOWN_SECS` | `60` (`DEFAULT_SCANNER_DEEP_VERIFY_COOLDOWN_SECS`, `scanner_folder.rs`) | Objects modified within this window are skipped by deep (bitrot) verification in the current cycle. |
|
||||||
| `RUSTFS_HEAL_OBJECT_SELECT_PROB` | `1024` (`DEFAULT_HEAL_OBJECT_SELECT_PROB`, `scanner_folder.rs`) | Sampling divisor for scanner-originated heal checks: roughly one object in N per cycle is selected for a low-priority heal check. |
|
| `RUSTFS_HEAL_OBJECT_SELECT_PROB` | `1024` (`DEFAULT_HEAL_OBJECT_SELECT_PROB`, `scanner_folder.rs`) | Sampling divisor for scanner-originated heal checks: roughly one object in N per cycle is selected for a low-priority heal check. `0` disables sampled checks. When N is smaller than the compacted-directory interval, every object in a selected directory is eligible; compaction must not round the sampling probability to zero. |
|
||||||
| `RUSTFS_DATA_USAGE_UPDATE_DIR_CYCLES` | `16` (`DATA_USAGE_UPDATE_DIR_CYCLES`, `scanner_folder.rs`) | Every N cycles a compacted directory is re-descended instead of reusing its cached usage. `1` forces re-descent every cycle (used by lifecycle e2e lanes). |
|
| `RUSTFS_DATA_USAGE_UPDATE_DIR_CYCLES` | `16` (`DATA_USAGE_UPDATE_DIR_CYCLES`, `scanner_folder.rs`) | Every N cycles a compacted directory is re-descended instead of reusing its cached usage. `1` forces re-descent every cycle (used by lifecycle e2e lanes); `0` is normalized to `1`. |
|
||||||
| `RUSTFS_DATA_USAGE_FAILED_OBJECT_TTL_SECS` | `86400` (`DEFAULT_FAILED_OBJECT_TTL_SECS`, `scanner_folder.rs`) | Retention of per-bucket failed-object entries in the usage cache. |
|
| `RUSTFS_DATA_USAGE_FAILED_OBJECT_TTL_SECS` | `86400` (`DEFAULT_FAILED_OBJECT_TTL_SECS`, `scanner_folder.rs`) | Retention of per-bucket failed-object retry entries in the usage cache. `0` disables and clears the retry cache; it does not allow failed scans to publish complete usage. Cached failures remain visible in each partial snapshot without extending their retry deadline. |
|
||||||
| `RUSTFS_DATA_USAGE_FAILED_OBJECTS_MAX` | `10000` (`DEFAULT_FAILED_OBJECTS_MAX`, `scanner_folder.rs`) | Cap on retained failed-object entries per bucket. |
|
| `RUSTFS_DATA_USAGE_FAILED_OBJECTS_MAX` | `10000` (`DEFAULT_FAILED_OBJECTS_MAX`, `scanner_folder.rs`) | Cap on retained failed-object entries per bucket. |
|
||||||
|
|
||||||
### Cycle budgets and cadence
|
### Cycle budgets and cadence
|
||||||
|
|||||||
Reference in New Issue
Block a user