From abc5f2e8188ad62416a0f7fed546fa8dbd8e58b4 Mon Sep 17 00:00:00 2001 From: Henry Guo Date: Thu, 30 Jul 2026 07:22:30 +0800 Subject: [PATCH] fix(scanner): persist portable usage cache keys (#5444) * fix(scanner): persist portable usage cache keys * fix(scanner): validate complete bucket cache graphs Co-Authored-By: heihutu --------- Co-authored-by: houseme Co-authored-by: heihutu --- .github/workflows/ci.yml | 4 + Cargo.lock | 7 - crates/data-usage/Cargo.toml | 1 - crates/data-usage/src/data_usage.rs | 53 ++- crates/ecstore/src/disk/local.rs | 21 +- crates/scanner/src/data_usage_define.rs | 434 +++++++++++++++++- crates/scanner/src/remote_scanner.rs | 64 +-- crates/scanner/src/scanner_folder.rs | 163 +++++-- crates/scanner/src/scanner_io.rs | 396 +++++++++++++--- .../tests/admin_diagnostic_capability_e2e.rs | 19 + 10 files changed, 1015 insertions(+), 147 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 27a806a2e..184771d4e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -173,6 +173,10 @@ jobs: run: cargo clippy --all-targets -- -D warnings - name: Run nextest tests + env: + # Three concurrent workspace test links saturate the self-hosted + # runner's overlay I/O and can wedge Cargo until the 75m timeout. + CARGO_BUILD_JOBS: "2" run: | mkdir -p artifacts/test-and-lint # Evidence sampler for issue #5394: the post-mortem pgrep below runs diff --git a/Cargo.lock b/Cargo.lock index 137b9505c..2cf81a4e3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7237,12 +7237,6 @@ dependencies = [ "path-dedot", ] -[[package]] -name = "path-clean" -version = "1.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "17359afc20d7ab31fdb42bb844c8b3bb1dabd7dcf7e68428492da7f16966fcef" - [[package]] name = "path-dedot" version = "4.0.1" @@ -9096,7 +9090,6 @@ name = "rustfs-data-usage" version = "1.0.0-beta.11" dependencies = [ "async-trait", - "path-clean", "rmp-serde", "rustfs-filemeta", "serde", diff --git a/crates/data-usage/Cargo.toml b/crates/data-usage/Cargo.toml index 790a0b977..b239e3dd6 100644 --- a/crates/data-usage/Cargo.toml +++ b/crates/data-usage/Cargo.toml @@ -29,7 +29,6 @@ workspace = true [dependencies] serde = { workspace = true, features = ["derive"] } -path-clean = { workspace = true } rmp-serde = { workspace = true } async-trait = { workspace = true } rustfs-filemeta = { workspace = true } diff --git a/crates/data-usage/src/data_usage.rs b/crates/data-usage/src/data_usage.rs index 89a4933d6..0c36ad1cf 100644 --- a/crates/data-usage/src/data_usage.rs +++ b/crates/data-usage/src/data_usage.rs @@ -12,12 +12,10 @@ // See the License for the specific language governing permissions and // limitations under the License. -use path_clean::PathClean; use serde::{Deserialize, Serialize}; use std::{ collections::{HashMap, HashSet}, hash::{DefaultHasher, Hash, Hasher}, - path::Path, time::{Duration, SystemTime}, }; @@ -1110,9 +1108,39 @@ fn mark(duc: &DataUsageCache, entry: &DataUsageEntry, found: &mut HashSet String { + let rooted = data.starts_with('/'); + let mut parts = Vec::new(); + + for part in data.split('/') { + match part { + "" | "." => {} + ".." => { + if parts.last().is_some_and(|last| *last != "..") { + parts.pop(); + } else if !rooted { + parts.push(part); + } + } + _ => parts.push(part), + } + } + + let clean = parts.join("/"); + match (rooted, clean.is_empty()) { + (true, true) => "/".to_string(), + (true, false) => format!("/{clean}"), + (false, true) => ".".to_string(), + (false, false) => clean, + } +} + +/// Hash a slash-separated path for data usage caching. +/// +/// Cache identifiers are persisted and exchanged across nodes, so their +/// normalization must not depend on the host operating system. pub fn hash_path(data: &str) -> DataUsageHash { - DataUsageHash(Path::new(&data).clean().to_string_lossy().to_string()) + DataUsageHash(clean_data_usage_path(data)) } impl DataUsageInfo { @@ -1497,6 +1525,23 @@ mod tests { buckets_count: u64, } + #[test] + fn hash_path_uses_portable_slash_semantics() { + for (input, expected) in [ + ("", "."), + (".", "."), + ("/", "/"), + ("//bucket///prefix/", "/bucket/prefix"), + ("bucket/./prefix//object", "bucket/prefix/object"), + ("bucket/a/../b", "bucket/b"), + ("../bucket/..", ".."), + ("/../../bucket", "/bucket"), + ("bucket\\prefix/object", "bucket\\prefix/object"), + ] { + assert_eq!(hash_path(input).key(), expected, "unexpected portable cache key for {input:?}"); + } + } + #[test] fn completeness_marker_is_additive_for_legacy_named_readers() { let current = DataUsageInfo { diff --git a/crates/ecstore/src/disk/local.rs b/crates/ecstore/src/disk/local.rs index 94b85b213..6dc8a829a 100644 --- a/crates/ecstore/src/disk/local.rs +++ b/crates/ecstore/src/disk/local.rs @@ -4555,9 +4555,11 @@ impl LocalDisk { let tmp_path = Self::meta_path(root, RUSTFS_META_TMP_BUCKET); let tmp_old_path = Self::meta_path(root, RUSTFS_META_TMP_OLD_BUCKET).join(Uuid::new_v4().to_string()); - rename_all(&tmp_path, &tmp_old_path, root).await.inspect_err(|err| { - log_startup_disk_error("cleanup_tmp_rename_all", &tmp_path, err); - })?; + rename_all_ignore_missing_source(&tmp_path, &tmp_old_path, root) + .await + .inspect_err(|err| { + log_startup_disk_error("cleanup_tmp_rename_all", &tmp_path, err); + })?; let tmp_deleted_path = Self::meta_path(root, RUSTFS_META_TMP_DELETED_BUCKET); tokio::fs::create_dir_all(&tmp_deleted_path).await.inspect_err(|err| { @@ -13065,6 +13067,19 @@ mod test { assert!(format_info.last_check.is_none(), "cached format timestamp should be cleared"); } + #[tokio::test] + async fn cleanup_tmp_on_startup_allows_missing_tmp_directory() { + use tempfile::tempdir; + + let dir = tempdir().expect("operation should succeed"); + + LocalDisk::cleanup_tmp_on_startup(dir.path(), Arc::new(AtomicU32::new(0)), Arc::new(Notify::new())) + .await + .expect("missing temporary directory should already be clean"); + + assert!(LocalDisk::meta_path(dir.path(), RUSTFS_META_TMP_DELETED_BUCKET).exists()); + } + #[tokio::test] async fn cleanup_tmp_on_startup_moves_existing_tmp_and_recreates_trash() { use tempfile::tempdir; diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 70efe5856..1d34cdc7a 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -48,6 +48,7 @@ pub const DATA_USAGE_ROOT: &str = SLASH_SEPARATOR; const DATA_USAGE_BLOOM_NAME: &str = ".bloomcycle.bin"; pub const DATA_USAGE_CACHE_NAME: &str = ".usage-cache.bin"; +pub(crate) const DATA_USAGE_CACHE_KEY_FORMAT: u16 = 1; const DATA_USAGE_CACHE_SAVE_RETRIES: u32 = 2; const DATA_USAGE_CACHE_BACKUP_SAVE_TIMEOUT_SECS_MAX: u64 = 5; const DATA_USAGE_CACHE_BACKUP_SAVE_RETRIES: u32 = 0; @@ -437,6 +438,8 @@ pub struct DataUsageCacheInfo { pub snapshot_complete: bool, #[serde(default)] pub scan_plan_digest: Option, + #[serde(default)] + pub cache_key_format: u16, } impl Serialize for DataUsageCacheInfo { @@ -446,7 +449,7 @@ impl Serialize for DataUsageCacheInfo { { // Keep this metadata map-encoded so older readers can ignore fields // appended by newer scanner versions during rolling upgrades. - let mut state = serializer.serialize_map(Some(15))?; + let mut state = serializer.serialize_map(Some(16))?; state.serialize_entry("name", &self.name)?; state.serialize_entry("next_cycle", &self.next_cycle)?; state.serialize_entry("leader_epoch", &self.leader_epoch)?; @@ -462,6 +465,7 @@ impl Serialize for DataUsageCacheInfo { state.serialize_entry("source", &self.source)?; state.serialize_entry("snapshot_complete", &self.snapshot_complete)?; state.serialize_entry("scan_plan_digest", &self.scan_plan_digest)?; + state.serialize_entry("cache_key_format", &self.cache_key_format)?; state.end() } } @@ -500,19 +504,34 @@ impl DataUsageCache { let source_matches = self.info.source == Some(source); let plan_matches = self.info.scan_plan_digest == Some(scan_plan_digest); - let reusable = self.info.name == name + let metadata_is_reusable = self.info.name == name && self.info.leader_epoch == leader_epoch && plan_matches - && (source_matches || (!require_source && self.info.source.is_none())); + && (source_matches || (!require_source && self.info.source.is_none())) + && self.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT; + let reusable = metadata_is_reusable + && (self.cache.is_empty() + || if name == DATA_USAGE_ROOT || self.info.snapshot_complete { + self.checked_flatten_complete_scope(name).is_some() + } else { + self.checked_flatten(name).is_some() + }); if !reusable { + let pending_heals = if self.info.name == name { + std::mem::take(&mut self.info.pending_heals) + } else { + Vec::new() + }; *self = Self::default(); self.info.name = name.to_string(); + self.info.pending_heals = pending_heals; } self.info.next_cycle = next_cycle; self.info.leader_epoch = leader_epoch; self.info.source = Some(source); self.info.scan_plan_digest = Some(scan_plan_digest); + self.info.cache_key_format = DATA_USAGE_CACHE_KEY_FORMAT; self.info.snapshot_complete = false; if reusable { DataUsageCachePrepareOutcome::Reused @@ -576,36 +595,60 @@ impl DataUsageCache { } pub(crate) fn checked_flatten(&self, path: &str) -> Option { + self.checked_flatten_inner(path).map(|(entry, _)| entry) + } + + pub(crate) fn checked_flatten_complete(&self, path: &str) -> Option { + self.checked_flatten_inner(path) + .filter(|(_, visited)| *visited == self.cache.len()) + .map(|(entry, _)| entry) + } + + pub(crate) fn checked_flatten_complete_scope(&self, path: &str) -> Option { + if path == DATA_USAGE_ROOT { + return self.checked_flatten_complete(path); + } + let (entry, visited) = self.checked_flatten_inner(path)?; + let root_parent_only = { + let path_key = hash_path(path).key(); + self.cache + .get(DATA_USAGE_ROOT) + .is_some_and(|root| root_is_parent_only(root, &path_key)) + }; + let expected_entries = self.cache.len().saturating_sub(usize::from(root_parent_only)); + (visited == expected_entries).then_some(entry) + } + + fn checked_flatten_inner(&self, path: &str) -> Option<(DataUsageEntry, usize)> { let root_key = hash_path(path).key(); - let root = self.cache.get(&root_key)?; - let mut visited = HashSet::from([root_key]); - let mut pending = root.children.iter().map(|child| (child.clone(), 1usize)).collect::>(); + let (root_key, root) = self.cache.get_key_value(&root_key)?; + if root.compacted && !root.children.is_empty() { + return None; + } + let mut visited = HashSet::from([root_key.as_str()]); + let mut pending = root.children.iter().map(|child| (child.as_str(), 1usize)).collect::>(); let mut flattened = DataUsageEntry::default(); - let mut root_entry = root.clone(); - root_entry.children.clear(); - if !flattened.checked_merge(&root_entry) { + if !flattened.checked_merge(root) { return None; } flattened.compacted = root.compacted; while let Some((key, depth)) = pending.pop() { - if depth > MAX_DATA_USAGE_CACHE_DEPTH || !visited.insert(key.clone()) { + if depth > MAX_DATA_USAGE_CACHE_DEPTH || !visited.insert(key) { return None; } - let entry = self.cache.get(&key)?; - if depth == MAX_DATA_USAGE_CACHE_DEPTH && !entry.children.is_empty() { + let entry = self.cache.get(key)?; + if (entry.compacted || depth == MAX_DATA_USAGE_CACHE_DEPTH) && !entry.children.is_empty() { return None; } - pending.extend(entry.children.iter().map(|child| (child.clone(), depth + 1))); + pending.extend(entry.children.iter().map(|child| (child.as_str(), depth + 1))); - let mut child_entry = entry.clone(); - child_entry.children.clear(); - if !flattened.checked_merge(&child_entry) { + if !flattened.checked_merge(entry) { return None; } } - Some(flattened) + Some((flattened, visited.len())) } fn flatten_with_guard(&self, root: &DataUsageEntry, visited: &mut HashSet, depth: usize) -> DataUsageEntry { @@ -1524,6 +1567,18 @@ fn mark_with_depth(duc: &DataUsageCache, entry: &DataUsageEntry, found: &mut Has } } +fn root_is_parent_only(root: &DataUsageEntry, child: &str) -> bool { + root.children.len() == 1 + && root.children.contains(child) + && root.size == 0 + && root.objects == 0 + && root.versions == 0 + && root.delete_markers == 0 + && root.replication_stats.is_none() + && !root.compacted + && root.failed_objects == 0 +} + /// Trait for storage-specific operations on DataUsageCache #[async_trait::async_trait] pub trait DataUsageCacheStorage { @@ -2287,6 +2342,7 @@ mod tests { assert!(decoded.source.is_none()); assert!(!decoded.snapshot_complete); assert!(decoded.scan_plan_digest.is_none()); + assert_eq!(decoded.cache_key_format, 0); } #[test] @@ -2328,6 +2384,7 @@ mod tests { assert!(decoded.source.is_none()); assert!(!decoded.snapshot_complete); assert!(decoded.scan_plan_digest.is_none()); + assert_eq!(decoded.cache_key_format, 0); } #[test] @@ -2363,6 +2420,7 @@ mod tests { source: Some(DataUsageCacheSource::new(1, 2)), snapshot_complete: true, scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, ..Default::default() @@ -2381,6 +2439,7 @@ mod tests { assert_eq!(current.info.source, Some(DataUsageCacheSource::new(1, 2))); assert!(current.info.snapshot_complete); assert_eq!(current.info.scan_plan_digest, Some(TEST_PLAN_DIGEST)); + assert_eq!(current.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); assert_eq!(current.find("bucket").map(|entry| entry.objects), Some(3)); let decoded: OldDataUsageCache = rmp_serde::from_slice(&buf).expect("Old reader failed to deserialize new cache"); @@ -2442,6 +2501,7 @@ mod tests { source: Some(source), snapshot_complete: false, scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, ..Default::default() @@ -2466,6 +2526,254 @@ mod tests { assert!(!cache.info.snapshot_complete); } + #[test] + fn data_usage_cache_prepare_for_scan_preserves_pending_heal_only_progress() { + let source = DataUsageCacheSource::new(1, 0); + let pending_heal = PendingScannerHeal { + kind: PendingScannerHealKind::Object, + bucket: "bucket".to_string(), + object: Some("prefix/object".to_string()), + version_id: Some("version-a".to_string()), + scan_mode: HealScanMode::Deep, + first_seen: 1, + last_attempt: 2, + attempts: 3, + last_admission_result: "full".to_string(), + last_admission_reason: "queue_full".to_string(), + }; + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + source: Some(source), + scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + pending_heals: vec![pending_heal.clone()], + ..Default::default() + }, + ..Default::default() + }; + + let outcome = cache.prepare_for_scan("bucket", 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reused); + assert_eq!(cache.info.pending_heals, vec![pending_heal]); + assert!(cache.cache.is_empty()); + assert!(!cache.info.snapshot_complete); + } + + #[test] + fn data_usage_cache_prepare_for_scan_preserves_namespace_pending_heals_during_key_format_rebuild() { + let source = DataUsageCacheSource::new(1, 0); + let pending_heal = PendingScannerHeal { + kind: PendingScannerHealKind::Object, + bucket: "bucket".to_string(), + object: Some("prefix/object".to_string()), + version_id: Some("version-a".to_string()), + scan_mode: HealScanMode::Deep, + first_seen: 1, + last_attempt: 2, + attempts: 3, + last_admission_result: "full".to_string(), + last_admission_reason: "queue_full".to_string(), + }; + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + next_cycle: 7, + source: Some(source), + scan_plan_digest: Some(TEST_PLAN_DIGEST), + pending_heals: vec![pending_heal.clone()], + ..Default::default() + }, + ..Default::default() + }; + cache.replace(DATA_USAGE_ROOT, "", DataUsageEntry::default()); + + let outcome = cache.prepare_for_scan(DATA_USAGE_ROOT, 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reset); + assert_eq!(cache.info.pending_heals, vec![pending_heal]); + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + + #[test] + fn data_usage_cache_prepare_for_scan_drops_pending_heals_from_a_different_scope() { + let source = DataUsageCacheSource::new(1, 0); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "old-bucket".to_string(), + next_cycle: 7, + source: Some(source), + scan_plan_digest: Some(TEST_PLAN_DIGEST), + pending_heals: vec![PendingScannerHeal { + kind: PendingScannerHealKind::Object, + bucket: "old-bucket".to_string(), + object: Some("prefix/object".to_string()), + version_id: None, + scan_mode: HealScanMode::Normal, + first_seen: 1, + last_attempt: 2, + attempts: 3, + last_admission_result: "full".to_string(), + last_admission_reason: "queue_full".to_string(), + }], + ..Default::default() + }, + ..Default::default() + }; + + let outcome = cache.prepare_for_scan("new-bucket", 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reset); + assert_eq!(cache.info.name, "new-bucket"); + assert!(cache.info.pending_heals.is_empty()); + } + + #[test] + fn data_usage_cache_prepare_for_scan_resets_unknown_key_format() { + let source = DataUsageCacheSource::new(1, 0); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + source: Some(source), + scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT + 1, + ..Default::default() + }, + ..Default::default() + }; + cache.replace( + "bucket", + "", + DataUsageEntry { + objects: 3, + ..Default::default() + }, + ); + + let outcome = cache.prepare_for_scan("bucket", 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reset); + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + + #[test] + fn data_usage_cache_prepare_for_scan_resets_persisted_windows_key_mismatch() { + let source = DataUsageCacheSource::new(1, 0); + let root_key = hash_path("bucket").key(); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + source: Some(source), + snapshot_complete: true, + scan_plan_digest: Some(TEST_PLAN_DIGEST), + ..Default::default() + }, + ..Default::default() + }; + cache.cache.insert( + root_key, + DataUsageEntry { + children: HashSet::from(["bucket/prefix".to_string()]), + ..Default::default() + }, + ); + cache.cache.insert( + "bucket\\prefix".to_string(), + DataUsageEntry { + objects: 3, + ..Default::default() + }, + ); + + let encoded = cache.marshal_msg().expect("legacy Windows cache should serialize"); + let mut decoded = DataUsageCache::unmarshal(&encoded).expect("legacy Windows cache should deserialize"); + let outcome = decoded.prepare_for_scan("bucket", 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reset); + assert!(decoded.cache.is_empty()); + assert_eq!(decoded.info.name, "bucket"); + assert_eq!(decoded.info.next_cycle, 8); + assert_eq!(decoded.info.source, Some(source)); + assert!(!decoded.info.snapshot_complete); + } + + #[test] + fn data_usage_cache_prepare_for_scan_resets_current_cache_with_dangling_child() { + let source = DataUsageCacheSource::new(1, 0); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + source: Some(source), + scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + ..Default::default() + }; + cache.cache.insert( + hash_path("bucket").key(), + DataUsageEntry { + children: HashSet::from([hash_path("bucket/missing").key()]), + ..Default::default() + }, + ); + + let outcome = cache.prepare_for_scan("bucket", 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reset); + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + + #[test] + fn data_usage_cache_prepare_for_scan_resets_complete_bucket_cache_with_detached_entry() { + let source = DataUsageCacheSource::new(1, 0); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + source: Some(source), + snapshot_complete: true, + scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + ..Default::default() + }; + cache.cache.insert( + hash_path("bucket").key(), + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + cache.cache.insert( + hash_path("bucket/detached").key(), + DataUsageEntry { + objects: 2, + ..Default::default() + }, + ); + assert_eq!(cache.checked_flatten("bucket").map(|entry| entry.objects), Some(1)); + assert!( + cache.checked_flatten_complete("bucket").is_none(), + "complete bucket cache reuse must reject detached entries" + ); + + let outcome = cache.prepare_for_scan("bucket", 8, 0, source, TEST_PLAN_DIGEST, true); + + assert_eq!(outcome, DataUsageCachePrepareOutcome::Reset); + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + #[test] fn data_usage_cache_prepare_for_scan_rejects_legacy_cache_without_a_bucket_plan() { let source = DataUsageCacheSource::new(0, 0); @@ -2836,6 +3144,98 @@ mod tests { ); } + #[test] + fn checked_flatten_complete_rejects_detached_entries() { + let mut cache = DataUsageCache::default(); + cache.cache.insert( + hash_path("bucket").key(), + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + cache.cache.insert( + hash_path("bucket/detached").key(), + DataUsageEntry { + objects: 2, + ..Default::default() + }, + ); + + assert_eq!( + cache.checked_flatten("bucket").map(|entry| entry.objects), + Some(1), + "subtree flattening may ignore entries outside the requested subtree" + ); + assert!( + cache.checked_flatten_complete("bucket").is_none(), + "an authoritative cache root must reach every persisted entry" + ); + } + + #[test] + fn checked_flatten_rejects_compacted_entries_with_children() { + let root_key = hash_path("bucket").key(); + let child_key = hash_path("bucket/prefix").key(); + let mut cache = DataUsageCache::default(); + cache.cache.insert( + root_key, + DataUsageEntry { + children: HashSet::from([child_key.clone()]), + compacted: true, + ..Default::default() + }, + ); + cache.cache.insert( + child_key, + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + + assert!( + cache.checked_flatten_complete("bucket").is_none(), + "a compacted entry cannot retain child links without double-counting" + ); + } + + #[test] + fn checked_flatten_rejects_compacted_descendants_with_children() { + let root_key = hash_path("bucket").key(); + let child_key = hash_path("bucket/prefix").key(); + let grandchild_key = hash_path("bucket/prefix/object").key(); + let mut cache = DataUsageCache::default(); + cache.cache.insert( + root_key, + DataUsageEntry { + children: HashSet::from([child_key.clone()]), + ..Default::default() + }, + ); + cache.cache.insert( + child_key, + DataUsageEntry { + objects: 1, + children: HashSet::from([grandchild_key.clone()]), + compacted: true, + ..Default::default() + }, + ); + cache.cache.insert( + grandchild_key, + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + + assert!( + cache.checked_flatten("bucket").is_none(), + "a compacted descendant cannot retain child links without double-counting" + ); + } + #[test] fn checked_flatten_accepts_depth_limit_and_rejects_deeper_tree() { let root_key = hash_path("bucket").key(); diff --git a/crates/scanner/src/remote_scanner.rs b/crates/scanner/src/remote_scanner.rs index 1164587d0..c38f6a8bb 100644 --- a/crates/scanner/src/remote_scanner.rs +++ b/crates/scanner/src/remote_scanner.rs @@ -14,8 +14,8 @@ use crate::scanner_budget::{ScannerCycleBudget, ScannerCycleBudgetConfig}; use crate::scanner_io::{ - ScannerDiskScanOutcome, ScannerIODisk, cache_root_entry_info, cache_snapshot_is_current, scanner_cache_lock_resource, - scanner_cache_lock_timeout, scanner_set_disk_inventory, + DataUsageCacheScanState, ScannerDiskScanOutcome, ScannerIODisk, cache_root_entry_info, current_cache_root_or_prepare, + scanner_cache_lock_resource, scanner_cache_lock_timeout, scanner_set_disk_inventory, }; use crate::storage_api::owner::NS_SCANNER_PROTOCOL_VERSION; use crate::storage_api::scan::NamespaceLocking as _; @@ -966,32 +966,34 @@ async fn scan_and_persist_local_bucket( let revisions = cache.load_with_revisions(set.clone(), &cache_name).await.map_err(|err| { RemoteScannerServerError::worker(format!("remote namespace scanner cache load or revision lookup failed: {err}")) })?; - if cache_snapshot_is_current(&cache, &bucket, source, next_cycle, leader_epoch, scan_plan_digest) { - if guard.is_lock_lost() { - return Err(RemoteScannerServerError::worker( - "remote namespace scanner cache lock was lost before reusing the current snapshot", - )); + let scan_state = current_cache_root_or_prepare(&mut cache, &bucket, source, next_cycle, leader_epoch, scan_plan_digest, true); + match scan_state { + DataUsageCacheScanState::Current(usage) => { + if guard.is_lock_lost() { + return Err(RemoteScannerServerError::worker( + "remote namespace scanner cache lock was lost before reusing the current snapshot", + )); + } + return Ok(RemoteScannerFrameResult::Complete(Box::new(RemoteScannerComplete { + source, + scan_plan_digest, + usage: *usage, + pending_maintenance_work: !cache.info.pending_heals.is_empty(), + }))); } - return Ok(RemoteScannerFrameResult::Complete(Box::new(RemoteScannerComplete { - source, - scan_plan_digest, - usage: cache_root_entry_info(&cache) - .map_err(|err| RemoteScannerServerError::worker(format!("remote namespace scanner cache is corrupt: {err}")))?, - pending_maintenance_work: !cache.info.pending_heals.is_empty(), - }))); - } - match cache.prepare_for_scan(&bucket, next_cycle, leader_epoch, source, scan_plan_digest, true) { - DataUsageCachePrepareOutcome::RejectedNewerCycle => { - return Ok(RemoteScannerFrameResult::CycleAhead { - required_cycle: cache.info.next_cycle, - }); - } - DataUsageCachePrepareOutcome::RejectedNewerLeader => { - return Err(RemoteScannerServerError::worker( - "remote namespace scanner rejected work from an older leader epoch", - )); - } - DataUsageCachePrepareOutcome::Reused | DataUsageCachePrepareOutcome::Reset => {} + DataUsageCacheScanState::Prepared { outcome, .. } => match outcome { + DataUsageCachePrepareOutcome::RejectedNewerCycle => { + return Ok(RemoteScannerFrameResult::CycleAhead { + required_cycle: cache.info.next_cycle, + }); + } + DataUsageCachePrepareOutcome::RejectedNewerLeader => { + return Err(RemoteScannerServerError::worker( + "remote namespace scanner rejected work from an older leader epoch", + )); + } + DataUsageCachePrepareOutcome::Reused | DataUsageCachePrepareOutcome::Reset => {} + }, } cache.info.skip_healing = skip_healing; @@ -1971,6 +1973,7 @@ mod tests { #[test] fn request_rejects_empty_truncated_oversized_and_wrong_version_payloads() { + assert_eq!(NS_SCANNER_PROTOCOL_VERSION, 3); assert!(decode_remote_scanner_request(&[]).is_err()); let mut body = rmp_serde::to_vec_named(&test_request(Uuid::new_v4())).expect("request should encode"); @@ -1980,7 +1983,7 @@ mod tests { let oversized = vec![0_u8; NS_SCANNER_MAX_REQUEST_BODY_SIZE + 1]; assert!(decode_remote_scanner_request(&oversized).is_err()); - for version in [NS_SCANNER_PROTOCOL_VERSION - 1, NS_SCANNER_PROTOCOL_VERSION + 1] { + for version in [2, 4] { let mut wrong_version = test_request(Uuid::new_v4()); wrong_version.version = version; let body = rmp_serde::to_vec_named(&wrong_version).expect("request should encode"); @@ -2885,14 +2888,15 @@ mod tests { #[tokio::test] async fn wrong_frame_version_and_sequence_are_rejected() { + assert_eq!(NS_SCANNER_PROTOCOL_VERSION, 3); let request_id = Uuid::new_v4(); let auth = FrameAuthenticator::for_test(request_id); let frame = RemoteScannerFrame::progress(RemoteScannerProgress::default()); let payload = rmp_serde::to_vec_named(&frame).expect("frame should encode"); for (version, sequence, expected_error) in [ - (NS_SCANNER_PROTOCOL_VERSION - 1, 0, "unsupported remote namespace scanner frame version"), - (NS_SCANNER_PROTOCOL_VERSION + 1, 0, "unsupported remote namespace scanner frame version"), + (2, 0, "unsupported remote namespace scanner frame version"), + (4, 0, "unsupported remote namespace scanner frame version"), (NS_SCANNER_PROTOCOL_VERSION, 1, "frame sequence is invalid"), ] { let envelope = RemoteScannerFrameEnvelope { diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index ad05c1cfc..7d2ea6608 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -1508,6 +1508,9 @@ impl FolderScanner { if entry.children.contains(&child) { continue; } + if !self.old_cache.cache.contains_key(&child) { + continue; + } let child_hash = DataUsageHash(child.clone()); self.new_cache @@ -2320,7 +2323,6 @@ impl FolderScanner { } tokio::task::yield_now().await; - let h = DataUsageHash(folder_item.name.clone()); into.add_child(&h); self.record_scan_resume_hint(&folder_item.name); // We scanned a folder, optionally send update. @@ -2646,7 +2648,7 @@ impl FolderScanner { tokio::task::yield_now().await; } else { let mut dst = DataUsageEntry::default(); - let h = DataUsageHash(folder_item.name.clone()); + let h = hash_path(&folder_item.name); // Use Box::pin for recursive async call let fut = Box::pin(self.scan_folder(ctx.clone(), folder_item.clone(), &mut dst)); @@ -2911,7 +2913,7 @@ mod tests { use serial_test::serial; #[cfg(unix)] use std::os::unix::fs::{PermissionsExt, symlink}; - use std::sync::atomic::AtomicBool; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use temp_env::{with_var, with_var_unset}; use uuid::Uuid; @@ -3864,11 +3866,15 @@ mod tests { } async fn write_test_object_metadata(root: &std::path::Path, bucket: &str, object: &str) { + write_test_object_metadata_bytes(root, bucket, object, &metadata_for_object(bucket, object)).await; + } + + async fn write_test_object_metadata_bytes(root: &std::path::Path, bucket: &str, object: &str, metadata: &[u8]) { let object_dir = root.join(bucket).join(object); tokio::fs::create_dir_all(&object_dir) .await .expect("failed to create test object directory"); - tokio::fs::write(object_dir.join("xl.meta"), metadata_for_object(bucket, object)) + tokio::fs::write(object_dir.join("xl.meta"), metadata) .await .expect("failed to write test object metadata"); } @@ -4165,27 +4171,28 @@ mod tests { async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() { let (mut scanner, temp_dir) = build_test_scanner().await; let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir.clone()); - let _heal_responder = rustfs_common::heal_channel::init_heal_channel().ok().map(|mut heal_rx| { - tokio::spawn(async move { - while let Some(command) = heal_rx.recv().await { - if let rustfs_common::heal_channel::HealChannelCommand::Start { response_tx, .. } = command { - let _ = response_tx.send(Ok(HealAdmissionResult::Accepted)); - } + let heal_starts = Arc::new(AtomicUsize::new(0)); + let heal_starts_clone = heal_starts.clone(); + let mut heal_rx = + rustfs_common::heal_channel::init_heal_channel().expect("heal channel should initialize once for scanner tests"); + let _heal_responder = tokio::spawn(async move { + while let Some(command) = heal_rx.recv().await { + if let rustfs_common::heal_channel::HealChannelCommand::Start { response_tx, .. } = command { + heal_starts_clone.fetch_add(1, Ordering::Relaxed); + let _ = response_tx.send(Ok(HealAdmissionResult::Accepted)); } - }) + } }); let bucket = "src-archive"; - tokio::fs::create_dir_all(temp_dir.join(bucket)) - .await - .expect("failed to create bucket directory"); + let object = "snapshots/37b3f20d941e2f5e6d99114d9bb2f3e67a8a2e5c9c4c5a1b0d6e7f8091a2b3c4"; + let metadata = metadata_for_object(bucket, object); + write_test_object_metadata_bytes(&temp_dir, bucket, object, &metadata).await; let mut disks = vec![scanner.local_disk.clone()]; for disk_name in ["disk2", "disk3", "disk4"] { let disk_root = temp_dir.join(disk_name); - tokio::fs::create_dir_all(disk_root.join(bucket)) - .await - .expect("failed to create extra disk bucket directory"); + write_test_object_metadata_bytes(&disk_root, bucket, object, &metadata).await; let endpoint = Endpoint::try_from(disk_root.to_string_lossy().as_ref()).expect("failed to create extra disk endpoint"); let disk = new_disk( @@ -4204,7 +4211,7 @@ mod tests { scanner.disks = disks; scanner.disks_quorum = 2; scanner.old_cache.replace( - "src-archive/snapshots/37b3f20d941e2f5e6d99114d9bb2f3e67a8a2e5c9c4c5a1b0d6e7f8091a2b3c4", + &format!("{bucket}/{object}"), bucket, DataUsageEntry { objects: 1, @@ -4219,13 +4226,17 @@ mod tests { object_heal_prob_div: 1, }; - tokio::time::timeout( - Duration::from_millis(200), - scanner.scan_folder(CancellationToken::new(), folder, &mut into), - ) - .await - .expect("scan_folder should not hang after list_path_raw finishes") - .expect("scan_folder should finish successfully"); + tokio::time::timeout(Duration::from_secs(2), scanner.scan_folder(CancellationToken::new(), folder, &mut into)) + .await + .expect("scan_folder should not hang after list_path_raw finishes") + .expect("scan_folder should finish successfully"); + + let root = scanner + .new_cache + .checked_flatten(bucket) + .expect("healed cache must contain canonical child links"); + assert_eq!(root.objects, 1); + assert!(heal_starts.load(Ordering::Relaxed) > 0, "test must execute the heal child-link path"); } #[tokio::test] @@ -4774,7 +4785,7 @@ mod tests { .expect("unbounded scan should finish after partial progress"); let root = result - .size_recursive("bucket") + .checked_flatten("bucket") .expect("completed cache should retain bucket usage"); assert_eq!(root.objects, 5); assert!(result.info.snapshot_complete); @@ -4828,6 +4839,106 @@ mod tests { ); } + #[tokio::test] + #[serial] + async fn test_partial_entry_does_not_carry_missing_old_child() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir), + }; + let root_hash = hash_path("bucket"); + scanner.old_cache.cache.insert( + root_hash.key(), + DataUsageEntry { + children: HashSet::from([hash_path("bucket/missing").key()]), + ..Default::default() + }, + ); + + let mut partial = DataUsageEntry { + objects: 2, + size: 2, + ..Default::default() + }; + scanner.carry_forward_old_children(&root_hash, &mut partial); + scanner.new_cache.replace_hashed(&root_hash, &None, &partial); + + assert!(partial.children.is_empty()); + let flattened = scanner + .new_cache + .checked_flatten("bucket") + .expect("a partial cache must not retain dangling child links"); + assert_eq!(flattened.objects, 2); + assert_eq!(flattened.size, 2); + } + + #[tokio::test] + #[serial] + async fn test_legacy_windows_cache_rebuilds_and_round_trips_portable_keys() { + let (scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir.clone()), + }; + write_test_object_metadata(&temp_dir, "bucket", "prefix/object").await; + + let source = crate::data_usage_define::DataUsageCacheSource::new(0, 0); + let scan_plan_digest = crate::data_usage_define::DataUsageScanPlanDigest([9; 32]); + let mut legacy = DataUsageCache { + info: crate::data_usage_define::DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 7, + source: Some(source), + scan_plan_digest: Some(scan_plan_digest), + ..Default::default() + }, + ..Default::default() + }; + legacy.cache.insert( + "bucket".to_string(), + DataUsageEntry { + children: HashSet::from(["bucket\\prefix".to_string()]), + ..Default::default() + }, + ); + legacy.cache.insert( + "bucket\\prefix".to_string(), + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + let encoded = legacy.marshal_msg().expect("legacy cache should serialize"); + let mut migrated = DataUsageCache::unmarshal(&encoded).expect("legacy cache should deserialize"); + assert_eq!( + migrated.prepare_for_scan("bucket", 8, 0, source, scan_plan_digest, true), + crate::data_usage_define::DataUsageCachePrepareOutcome::Reset + ); + + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&parent, Default::default()); + let rebuilt = scan_data_folder( + budget.token(), + budget, + vec![scanner.local_disk.clone()], + scanner.local_disk, + migrated, + None, + HealScanMode::Normal, + SCANNER_SLEEPER.clone(), + ) + .await + .expect("portable cache rebuild should complete"); + let persisted = rebuilt.marshal_msg().expect("rebuilt cache should serialize"); + let decoded = DataUsageCache::unmarshal(&persisted).expect("rebuilt cache should deserialize"); + + assert_eq!(decoded.info.cache_key_format, crate::data_usage_define::DATA_USAGE_CACHE_KEY_FORMAT); + assert!(decoded.cache.keys().all(|key| !key.contains('\\'))); + let root = decoded + .checked_flatten("bucket") + .expect("rebuilt persisted cache should have a complete root"); + assert_eq!(root.objects, 1); + } + #[tokio::test] #[serial] async fn test_scan_data_folder_success_clears_resume_hint() { diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index 591182e86..f44b77031 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +use crate::data_usage_define::DATA_USAGE_CACHE_KEY_FORMAT; use crate::scanner_budget::ScannerCycleBudget; use crate::scanner_folder::{ScannerItem, scan_data_folder}; use crate::sleeper::SCANNER_SLEEPER; @@ -830,7 +831,7 @@ pub(crate) fn cache_root_entry_info(cache: &DataUsageCache) -> std::result::Resu return Err(ScannerError::Other("scanner cache root name is empty".to_string())); } let entry = cache - .checked_flatten(&cache.info.name) + .checked_flatten_complete_scope(&cache.info.name) .ok_or_else(|| ScannerError::Other(format!("scanner cache root is missing or corrupt: {}", cache.info.name)))?; Ok(DataUsageEntryInfo { @@ -852,7 +853,7 @@ fn should_publish_completed_snapshot(completed_count: usize, total_count: usize, #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum NamespaceScannerWorkerMode { Coordinator, - RemoteV3(uuid::Uuid), + RemoteV4(uuid::Uuid), } fn namespace_scanner_workers( @@ -868,7 +869,7 @@ fn namespace_scanner_workers( workers.extend( remote_disks .into_iter() - .map(|(disk, server_epoch)| (disk, NamespaceScannerWorkerMode::RemoteV3(server_epoch))), + .map(|(disk, server_epoch)| (disk, NamespaceScannerWorkerMode::RemoteV4(server_epoch))), ); workers } @@ -958,7 +959,57 @@ where let _ = tokio::time::timeout(SCANNER_CACHE_LOCK_LOSS_SHUTDOWN_TIMEOUT, scan).await; } -pub(crate) fn cache_snapshot_is_current( +pub(crate) fn current_cache_root_entry( + cache: &DataUsageCache, + name: &str, + source: DataUsageCacheSource, + next_cycle: u64, + leader_epoch: u64, + scan_plan_digest: DataUsageScanPlanDigest, +) -> std::result::Result, ScannerError> { + let metadata_is_current = cache.info.name == name + && cache.info.source == Some(source) + && cache.info.snapshot_complete + && cache.info.scan_plan_digest == Some(scan_plan_digest) + && cache.info.last_update.is_some() + && cache.info.next_cycle == next_cycle + && cache.info.leader_epoch == leader_epoch + && cache.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT; + if !metadata_is_current { + return Ok(None); + } + + cache_root_entry_info(cache).map(Some) +} + +pub(crate) enum DataUsageCacheScanState { + Current(Box), + Prepared { + outcome: DataUsageCachePrepareOutcome, + invalid_current: Option, + }, +} + +pub(crate) fn current_cache_root_or_prepare( + cache: &mut DataUsageCache, + name: &str, + source: DataUsageCacheSource, + next_cycle: u64, + leader_epoch: u64, + scan_plan_digest: DataUsageScanPlanDigest, + require_source: bool, +) -> DataUsageCacheScanState { + match current_cache_root_entry(cache, name, source, next_cycle, leader_epoch, scan_plan_digest) { + Ok(Some(root)) => DataUsageCacheScanState::Current(Box::new(root)), + current => DataUsageCacheScanState::Prepared { + invalid_current: current.err(), + outcome: cache.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, require_source), + }, + } +} + +#[cfg(test)] +fn cache_snapshot_is_current( cache: &DataUsageCache, name: &str, source: DataUsageCacheSource, @@ -966,13 +1017,10 @@ pub(crate) fn cache_snapshot_is_current( leader_epoch: u64, scan_plan_digest: DataUsageScanPlanDigest, ) -> bool { - cache.info.name == name - && cache.info.source == Some(source) - && cache.info.snapshot_complete - && cache.info.scan_plan_digest == Some(scan_plan_digest) - && cache.info.last_update.is_some() - && cache.info.next_cycle == next_cycle - && cache.info.leader_epoch == leader_epoch + matches!( + current_cache_root_entry(cache, name, source, next_cycle, leader_epoch, scan_plan_digest), + Ok(Some(_)) + ) } fn completed_data_usage_info( @@ -1061,6 +1109,7 @@ mod publish_gate_tests { source: Some(source), snapshot_complete: false, scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, ..Default::default() @@ -1102,6 +1151,7 @@ mod publish_gate_tests { source: Some(source), snapshot_complete: true, scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, ..Default::default() @@ -1423,17 +1473,145 @@ mod publish_gate_tests { assert!(!cache_snapshot_is_current(&cache, DATA_USAGE_ROOT, source, 10, 2, TEST_PLAN_DIGEST)); } + #[test] + fn current_cache_snapshot_rejects_persisted_windows_key_mismatch() { + let source = DataUsageCacheSource::new(1, 2); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 10, + last_update: Some(SystemTime::UNIX_EPOCH), + source: Some(source), + snapshot_complete: true, + scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + ..Default::default() + }; + cache.cache.insert( + "bucket".to_string(), + DataUsageEntry { + children: HashSet::from(["bucket/prefix".to_string()]), + ..Default::default() + }, + ); + cache.cache.insert( + "bucket\\prefix".to_string(), + DataUsageEntry { + objects: 3, + ..Default::default() + }, + ); + + assert!(!cache_snapshot_is_current(&cache, "bucket", source, 10, 0, TEST_PLAN_DIGEST)); + match current_cache_root_or_prepare(&mut cache, "bucket", source, 10, 0, TEST_PLAN_DIGEST, true) { + DataUsageCacheScanState::Prepared { + outcome: DataUsageCachePrepareOutcome::Reset, + invalid_current: Some(_), + } => {} + _ => panic!("an invalid current cache must enter the rebuild path"), + } + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + + #[test] + fn current_cache_snapshot_rejects_structurally_valid_legacy_key_format() { + let source = DataUsageCacheSource::new(1, 2); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 10, + last_update: Some(SystemTime::UNIX_EPOCH), + source: Some(source), + snapshot_complete: true, + scan_plan_digest: Some(TEST_PLAN_DIGEST), + ..Default::default() + }, + ..Default::default() + }; + cache.cache.insert( + "bucket".to_string(), + DataUsageEntry { + children: HashSet::from(["bucket\\prefix".to_string()]), + ..Default::default() + }, + ); + cache.cache.insert( + "bucket\\prefix".to_string(), + DataUsageEntry { + objects: 3, + ..Default::default() + }, + ); + assert_eq!(cache.checked_flatten("bucket").map(|entry| entry.objects), Some(3)); + + match current_cache_root_or_prepare(&mut cache, "bucket", source, 10, 0, TEST_PLAN_DIGEST, true) { + DataUsageCacheScanState::Prepared { + outcome: DataUsageCachePrepareOutcome::Reset, + invalid_current: None, + } => {} + _ => panic!("a legacy key format must enter the rebuild path"), + } + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + + #[test] + fn current_cache_snapshot_rejects_current_bucket_cache_with_detached_entry() { + let source = DataUsageCacheSource::new(1, 2); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + next_cycle: 10, + last_update: Some(SystemTime::UNIX_EPOCH), + source: Some(source), + snapshot_complete: true, + scan_plan_digest: Some(TEST_PLAN_DIGEST), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, + ..Default::default() + }, + ..Default::default() + }; + cache.cache.insert( + "bucket".to_string(), + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + cache.cache.insert( + "bucket/detached".to_string(), + DataUsageEntry { + objects: 2, + ..Default::default() + }, + ); + assert_eq!(cache.checked_flatten("bucket").map(|entry| entry.objects), Some(1)); + + match current_cache_root_or_prepare(&mut cache, "bucket", source, 10, 0, TEST_PLAN_DIGEST, true) { + DataUsageCacheScanState::Prepared { + outcome: DataUsageCachePrepareOutcome::Reset, + invalid_current: Some(_), + } => {} + _ => panic!("a detached complete bucket cache must enter the rebuild path"), + } + assert!(cache.cache.is_empty()); + assert_eq!(cache.info.cache_key_format, DATA_USAGE_CACHE_KEY_FORMAT); + } + #[test] fn namespace_scanner_worker_selection_keeps_coordinator_fallback_disks() { let server_epoch = uuid::Uuid::new_v4(); - let workers = namespace_scanner_workers(vec!["local", "legacy-remote"], vec![("v3", server_epoch)]); + let workers = namespace_scanner_workers(vec!["local", "legacy-remote"], vec![("v4", server_epoch)]); assert_eq!( workers, vec![ ("local", NamespaceScannerWorkerMode::Coordinator), ("legacy-remote", NamespaceScannerWorkerMode::Coordinator), - ("v3", NamespaceScannerWorkerMode::RemoteV3(server_epoch)), + ("v4", NamespaceScannerWorkerMode::RemoteV4(server_epoch)), ] ); assert!(namespace_scanner_workers::<()>(Vec::new(), Vec::new()).is_empty()); @@ -1507,6 +1685,7 @@ mod publish_gate_tests { source: Some(source), snapshot_complete: true, scan_plan_digest: Some(first), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, ..Default::default() @@ -1595,6 +1774,15 @@ async fn send_cache_root_entry_info( pending_maintenance_work: &AtomicBool, ) -> std::result::Result<(), ScannerError> { let root = cache_root_entry_info(cache)?; + send_cache_root_entry(bucket_result_tx, root, cache, pending_maintenance_work).await +} + +async fn send_cache_root_entry( + bucket_result_tx: &mpsc::Sender, + root: DataUsageEntryInfo, + cache: &DataUsageCache, + pending_maintenance_work: &AtomicBool, +) -> std::result::Result<(), ScannerError> { record_bucket_pending_maintenance_work(cache, pending_maintenance_work); bucket_result_tx .send(root) @@ -1690,13 +1878,16 @@ async fn persist_and_publish_cache_snapshot( ); return None; } - if cache_snapshot_is_current( - &persisted, - DATA_USAGE_ROOT, - source, - cache_snapshot.info.next_cycle, - cache_snapshot.info.leader_epoch, - scan_plan_digest, + if matches!( + current_cache_root_entry( + &persisted, + DATA_USAGE_ROOT, + source, + cache_snapshot.info.next_cycle, + cache_snapshot.info.leader_epoch, + scan_plan_digest, + ), + Ok(Some(_)) ) { cache_snapshot = persisted; } else { @@ -2329,6 +2520,7 @@ impl ScannerIOCache for SetDisks { source: Some(source), snapshot_complete: true, scan_plan_digest: Some(scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, cache: HashMap::new(), @@ -2450,7 +2642,7 @@ impl ScannerIOCache for SetDisks { subsystem = LOG_SUBSYSTEM_IO, pool = self.pool_index, set = self.set_index, - v3_disks = remote_disk_count, + v4_disks = remote_disk_count, unsupported_remote_disks, state = "unsupported_remote_disks_using_coordinator", "Scanner set assigned remote disks without namespace scanner support to coordinator-driven workers" @@ -2547,6 +2739,7 @@ impl ScannerIOCache for SetDisks { source: Some(source), snapshot_complete: false, scan_plan_digest: Some(scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, cache: HashMap::new(), @@ -2638,7 +2831,7 @@ impl ScannerIOCache for SetDisks { let dirty_usage_buckets_clone = dirty_usage_buckets.clone(); let cache_cycle_floor_clone = cache_cycle_floor.clone(); let remote_server_epoch = match worker_mode { - NamespaceScannerWorkerMode::RemoteV3(server_epoch) => Some(server_epoch), + NamespaceScannerWorkerMode::RemoteV4(server_epoch) => Some(server_epoch), NamespaceScannerWorkerMode::Coordinator => None, }; futs.push(tokio::spawn(async move { @@ -2948,48 +3141,71 @@ impl ScannerIOCache for SetDisks { continue; } }; - if cache_snapshot_is_current(&cache, &bucket.name, source, want_cycle, leader_epoch, bucket_scan_plan_digest) - { - if cache_guard.is_lock_lost() { - record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; - error!( - target: "rustfs::scanner::io", - event = EVENT_SCANNER_CACHE_PERSIST_STATE, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_IO, - bucket = %bucket.name, - cache_name = %cache_name, - state = "lock_lost_before_reuse", - "Current scanner bucket cache root publish skipped after lock loss" - ); - continue; - } - if let Err(e) = - send_cache_root_entry_info(&bucket_result_tx_clone, &cache, &pending_maintenance_work_clone).await - { - record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; - error!( - target: "rustfs::scanner::io", - event = EVENT_SCANNER_DATA_USAGE_STREAM, - component = LOG_COMPONENT_SCANNER, - subsystem = LOG_SUBSYSTEM_IO, - bucket = %bucket.name, - state = "send_current_root_failed", - error = %e, - "Current scanner bucket cache root entry publish failed" - ); - } - continue; - } - - match cache.prepare_for_scan( + let scan_state = current_cache_root_or_prepare( + &mut cache, &bucket.name, + source, want_cycle, leader_epoch, - source, bucket_scan_plan_digest, require_cache_source, - ) { + ); + let outcome = match scan_state { + DataUsageCacheScanState::Current(root) => { + if cache_guard.is_lock_lost() { + record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; + error!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_CACHE_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + bucket = %bucket.name, + cache_name = %cache_name, + state = "lock_lost_before_reuse", + "Current scanner bucket cache root publish skipped after lock loss" + ); + continue; + } + if let Err(e) = + send_cache_root_entry(&bucket_result_tx_clone, *root, &cache, &pending_maintenance_work_clone) + .await + { + record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; + error!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_DATA_USAGE_STREAM, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + bucket = %bucket.name, + state = "send_current_root_failed", + error = %e, + "Current scanner bucket cache root entry publish failed" + ); + } + continue; + } + DataUsageCacheScanState::Prepared { + outcome, + invalid_current, + } => { + if let Some(e) = invalid_current { + warn!( + target: "rustfs::scanner::io", + event = EVENT_SCANNER_CACHE_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_IO, + bucket = %bucket.name, + cache_name = %cache_name, + state = "current_cache_invalid", + error = %e, + "Current scanner bucket cache is invalid; rebuilding" + ); + } + outcome + } + }; + + match outcome { DataUsageCachePrepareOutcome::RejectedNewerCycle => { cache_cycle_floor_clone.fetch_max(cache.info.next_cycle, Ordering::AcqRel); record_failed_dirty_bucket(&failed_dirty_buckets_clone, &bucket.name).await; @@ -3361,6 +3577,7 @@ impl ScannerIOCache for SetDisks { source: Some(source), snapshot_complete: false, scan_plan_digest: Some(scan_plan_digest), + cache_key_format: DATA_USAGE_CACHE_KEY_FORMAT, ..Default::default() }, cache: HashMap::new(), @@ -4668,6 +4885,67 @@ mod tests { root.add_child(&crate::hash_path("bucket/missing")); dangling.replace("bucket", DATA_USAGE_ROOT, root); assert!(cache_root_entry_info(&dangling).is_err()); + + let mut detached = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + ..Default::default() + }, + ..Default::default() + }; + detached.replace("bucket", DATA_USAGE_ROOT, DataUsageEntry::default()); + detached.replace( + "bucket/detached", + "", + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + assert!(cache_root_entry_info(&detached).is_err()); + + let mut detached_bucket = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + ..Default::default() + }, + ..Default::default() + }; + detached_bucket.replace( + "bucket", + DATA_USAGE_ROOT, + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + detached_bucket.replace( + "bucket/detached", + "", + DataUsageEntry { + objects: 1, + ..Default::default() + }, + ); + assert!(cache_root_entry_info(&detached_bucket).is_err()); + + let mut compacted_with_child = DataUsageCache { + info: DataUsageCacheInfo { + name: "bucket".to_string(), + ..Default::default() + }, + ..Default::default() + }; + compacted_with_child.replace( + "bucket", + DATA_USAGE_ROOT, + DataUsageEntry { + compacted: true, + ..Default::default() + }, + ); + compacted_with_child.replace("bucket/prefix", "bucket", DataUsageEntry::default()); + assert!(cache_root_entry_info(&compacted_with_child).is_err()); } #[test] diff --git a/rustfs/tests/admin_diagnostic_capability_e2e.rs b/rustfs/tests/admin_diagnostic_capability_e2e.rs index f639802b1..ddd3532b6 100644 --- a/rustfs/tests/admin_diagnostic_capability_e2e.rs +++ b/rustfs/tests/admin_diagnostic_capability_e2e.rs @@ -84,6 +84,22 @@ async fn signed_bytes_request( .expect("signed admin request") } +async fn wait_for_ready(client: &Client, endpoint: &str) { + let ready_url = format!("{endpoint}/health/ready"); + tokio::time::timeout(Duration::from_secs(10), async { + loop { + if let Ok(response) = client.get(&ready_url).send().await + && response.status().is_success() + { + return; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } + }) + .await + .expect("embedded server should become ready"); +} + #[tokio::test] async fn diagnostic_handlers_enforce_advertised_runtime_contract() { let port = match find_available_port() { @@ -100,9 +116,11 @@ async fn diagnostic_handlers_enforce_advertised_runtime_contract() { .expect("start embedded server"); let endpoint = server.endpoint(); let client = Client::builder() + .no_proxy() .timeout(Duration::from_secs(5)) .build() .expect("HTTP client"); + wait_for_ready(&client, &endpoint).await; let response = signed_bytes_request( &client, @@ -144,6 +162,7 @@ async fn diagnostic_handlers_enforce_advertised_runtime_contract() { ); let stalled_client = Client::builder() + .no_proxy() .timeout(Duration::from_secs(60)) .build() .expect("stalled HTTP client");