diff --git a/crates/scanner/src/raw_page_index.rs b/crates/scanner/src/raw_page_index.rs index e861bd5da..023f9f585 100644 --- a/crates/scanner/src/raw_page_index.rs +++ b/crates/scanner/src/raw_page_index.rs @@ -16,6 +16,14 @@ use serde::{Deserialize, Serialize}; use sha2::{Digest as _, Sha256}; use thiserror::Error; +mod writer; +pub(crate) use writer::RawEnumerationPageWriter; + +#[cfg(test)] +thread_local! { + pub(crate) static RAW_PAGE_DIGEST_ENTRIES: std::cell::Cell = const { std::cell::Cell::new(0) }; +} + const RAW_PAGE_INDEX_VERSION: u16 = 1; const RAW_PAGE_ENTRY_MAX_BYTES: usize = 16 * 1024; @@ -145,6 +153,13 @@ impl RawEnumerationPageIndex { } } + pub(crate) fn parent(&self) -> Option<&str> { + match &self.state { + RawEnumerationPageIndexState::Unsupported => None, + RawEnumerationPageIndexState::Supported(inner) => Some(&inner.parent), + } + } + pub fn status(&self) -> RawEnumerationPageOwnerStatus { match &self.state { RawEnumerationPageIndexState::Unsupported => RawEnumerationPageOwnerStatus::Unsupported, @@ -295,25 +310,7 @@ impl RawEnumerationPageIndex { return Err(RawEnumerationPageIndexError::StaleGeneration); } let entries_start = inner.validated_committed_entries()?.len(); - let Some(building) = inner.building.as_ref() else { - return Err(RawEnumerationPageIndexError::EmptyCommit); - }; - if building.entries.is_empty() { - return Err(RawEnumerationPageIndexError::EmptyCommit); - } - building.validate( - u64::try_from(inner.pages.len()).unwrap_or(u64::MAX), - u64::try_from(entries_start).unwrap_or(u64::MAX), - inner.page_entry_limit, - )?; - let Some(building) = inner.building.take() else { - return Err(RawEnumerationPageIndexError::EmptyCommit); - }; - let page = RawEnumerationPage::new(&inner.parent, building); - inner.complete = page.terminal; - inner.pages.push(page.clone()); - inner.generation = inner.generation.saturating_add(1); - Ok(page) + inner.commit_building_page(entries_start) } pub fn page(&self, page_index: u64, expected_digest: [u8; 32]) -> Result<&RawEnumerationPage, RawEnumerationPageIndexError> { @@ -349,6 +346,30 @@ impl RawEnumerationPageIndex { } impl RawEnumerationPageIndexInner { + // Callers validate the committed prefix at the restore/API boundary. The + // live writer maintains its entry count without revisiting that prefix. + fn commit_building_page(&mut self, entries_start: usize) -> Result { + let Some(building) = self.building.as_ref() else { + return Err(RawEnumerationPageIndexError::EmptyCommit); + }; + if building.entries.is_empty() { + return Err(RawEnumerationPageIndexError::EmptyCommit); + } + building.validate( + u64::try_from(self.pages.len()).unwrap_or(u64::MAX), + u64::try_from(entries_start).unwrap_or(u64::MAX), + self.page_entry_limit, + )?; + let Some(building) = self.building.take() else { + return Err(RawEnumerationPageIndexError::EmptyCommit); + }; + let page = RawEnumerationPage::new(&self.parent, building); + self.complete = page.terminal; + self.pages.push(page.clone()); + self.generation = self.generation.saturating_add(1); + Ok(page) + } + fn status(&self) -> RawEnumerationPageOwnerStatus { let indexed_entries = u64::try_from(self.indexed_entries()).unwrap_or(u64::MAX); if let Some(building) = &self.building { @@ -528,6 +549,8 @@ fn entry_sets_match(left: &[String], right: &[String]) -> bool { } fn raw_page_digest(parent: &str, building: &RawEnumerationPageBuilder) -> [u8; 32] { + #[cfg(test)] + RAW_PAGE_DIGEST_ENTRIES.with(|count| count.set(count.get().saturating_add(building.entries.len()))); let mut digest = Sha256::new(); update_digest(&mut digest, b"version", &RAW_PAGE_INDEX_VERSION.to_le_bytes()); update_digest(&mut digest, b"parent", parent.as_bytes()); diff --git a/crates/scanner/src/raw_page_index/writer.rs b/crates/scanner/src/raw_page_index/writer.rs new file mode 100644 index 000000000..9cd7f04da --- /dev/null +++ b/crates/scanner/src/raw_page_index/writer.rs @@ -0,0 +1,149 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::{ + RawEnumerationPageBuilder, RawEnumerationPageIndex, RawEnumerationPageIndexError, RawEnumerationPageIndexInner, + RawEnumerationPageIndexState, owner_entry_is_valid, +}; +use std::collections::{BTreeSet, HashSet}; + +/// Exclusive, validated owner of an index during one directory enumeration. +/// The auxiliary sets are never serialized: restored indexes must cross the +/// validation boundary again before any pages can be trusted or extended. +pub(crate) struct RawEnumerationPageWriter { + inner: RawEnumerationPageIndexInner, + indexed: HashSet, + pending: BTreeSet, + observed_complete: HashSet, + complete_source_changed: bool, + observations: usize, + revalidate_after_entries: usize, + committed_entries: usize, +} + +impl RawEnumerationPageWriter { + pub(crate) fn new(index: RawEnumerationPageIndex) -> Result { + let RawEnumerationPageIndexState::Supported(inner) = index.state else { + return Err(RawEnumerationPageIndexError::Unsupported); + }; + if inner.parent.is_empty() || inner.page_entry_limit == 0 { + return Err(RawEnumerationPageIndexError::CorruptIndex); + } + let entries = inner.validated_indexed_entries()?; + let entry_count = entries.len(); + let indexed: HashSet<_> = entries.into_iter().collect(); + if indexed.len() != entry_count + || inner.pages.iter().any(|page| page.entries.len() > inner.page_entry_limit) + || (inner.complete && inner.building.is_some()) + || inner.pages.last().is_some_and(|page| page.terminal != inner.complete) + { + return Err(RawEnumerationPageIndexError::CorruptIndex); + } + let committed_entries = entry_count.saturating_sub(inner.building.as_ref().map_or(0, |page| page.entries.len())); + Ok(Self { + inner, + indexed, + pending: BTreeSet::new(), + observed_complete: HashSet::new(), + complete_source_changed: false, + observations: 0, + revalidate_after_entries: entry_count, + committed_entries, + }) + } + + pub(crate) fn indexed_entry_count(&self) -> usize { + self.indexed.len() + } + + pub(crate) fn record_entry(&mut self, entry: &str) -> Result<(), RawEnumerationPageIndexError> { + if !owner_entry_is_valid(entry) { + return Err(RawEnumerationPageIndexError::InvalidEntry); + } + self.observations = self.observations.saturating_add(1); + if self.inner.complete || self.inner.building.as_ref().is_some_and(|page| page.terminal) { + if self.indexed.contains(entry) { + self.observed_complete.insert(entry.to_owned()); + } else { + self.complete_source_changed = true; + } + } + if self.inner.complete { + if self.observations >= self.revalidate_after_entries + && (self.observed_complete.len() != self.indexed.len() || self.complete_source_changed) + { + return Err(RawEnumerationPageIndexError::IdentityMismatch); + } + return Ok(()); + } + if !self.indexed.contains(entry) { + self.pending.insert(entry.to_owned()); + } + // A restored directory can be enumerated in a different order. Keep + // the old observation floor and append at most one new entry per call, + // in the same order as normalization of the cumulative source. + if self.observations < self.revalidate_after_entries { + return Ok(()); + } + if self + .inner + .building + .as_ref() + .is_some_and(|page| page.entries.len() >= self.inner.page_entry_limit) + { + return self.commit_page(); + } + if let Some(entry) = self.pending.pop_first() { + let building = self.inner.building.get_or_insert_with(|| RawEnumerationPageBuilder { + page_index: u64::try_from(self.inner.pages.len()).unwrap_or(u64::MAX), + entries_start: u64::try_from(self.committed_entries).unwrap_or(u64::MAX), + entries: Vec::new(), + terminal: false, + }); + self.indexed.insert(entry.clone()); + building.entries.push(entry); + building.entries.sort(); + building.terminal = false; + self.inner.generation = self.inner.generation.saturating_add(1); + } + if self + .inner + .building + .as_ref() + .is_some_and(|page| page.terminal || page.entries.len() >= self.inner.page_entry_limit) + { + self.commit_page()?; + } + Ok(()) + } + + fn commit_page(&mut self) -> Result<(), RawEnumerationPageIndexError> { + self.inner.commit_building_page(self.committed_entries)?; + self.committed_entries = self.indexed.len(); + Ok(()) + } + + pub(crate) fn checkpoint(&self) -> Result { + let mut inner = self.inner.clone(); + if inner.building.is_some() { + inner.commit_building_page(self.committed_entries)?; + } + Ok(RawEnumerationPageIndex { + state: RawEnumerationPageIndexState::Supported(inner), + }) + } +} + +#[cfg(test)] +mod tests; diff --git a/crates/scanner/src/raw_page_index/writer/tests.rs b/crates/scanner/src/raw_page_index/writer/tests.rs new file mode 100644 index 000000000..eb46fbd52 --- /dev/null +++ b/crates/scanner/src/raw_page_index/writer/tests.rs @@ -0,0 +1,223 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::raw_page_index::RAW_PAGE_DIGEST_ENTRIES; + +fn checkpoint(mut index: RawEnumerationPageIndex) -> RawEnumerationPageIndex { + if let RawEnumerationPageIndexState::Supported(inner) = &index.state + && inner.building.is_some() + { + index + .commit_building_page(index.generation().expect("supported index generation")) + .expect("reference checkpoint must commit its partial page"); + } + index +} + +fn compare_with_cumulative_ingestion(initial: RawEnumerationPageIndex, source: &[&str]) { + let mut reference = initial.clone(); + let floor = reference.indexed_entries().expect("reference entries").len(); + let mut writer = RawEnumerationPageWriter::new(initial).expect("validated writer"); + let mut observed = Vec::new(); + for entry in source { + observed.push((*entry).to_owned()); + let expected = if observed.len() < floor { + Ok(()) + } else { + reference + .ingest_partial_owner_entries(observed.clone(), 1, reference.generation().expect("reference generation")) + .and_then(|outcome| { + if outcome.ready_to_commit { + reference.commit_building_page(reference.generation().expect("commit generation"))?; + } + Ok(()) + }) + }; + let actual = writer.record_entry(entry); + assert_eq!(actual, expected, "observation {observed:?}"); + if actual.is_err() { + return; + } + let expected = checkpoint(reference.clone()); + let actual = writer.checkpoint().expect("incremental checkpoint"); + assert_eq!(actual, expected, "page identity and generation at {observed:?}"); + assert_eq!( + rmp_serde::to_vec(&actual).expect("encode incremental index"), + rmp_serde::to_vec(&expected).expect("encode reference index"), + "persisted format must remain unchanged" + ); + assert_eq!(writer.indexed_entry_count(), actual.indexed_entries().expect("validate output").len()); + } +} + +#[test] +fn raw_enumeration_writer_matches_cumulative_pages_with_duplicates_and_restarts() { + for limit in [1, 2, 3, 128] { + let source = ["c", "a", "b", "b", "f", "d", "e", "g"]; + let empty = RawEnumerationPageIndex::new("bucket/metadata", limit).expect("empty index"); + compare_with_cumulative_ingestion(empty.clone(), &source); + let mut writer = RawEnumerationPageWriter::new(empty).expect("writer"); + for entry in source { + writer.record_entry(entry).expect("entry before interruption"); + let saved = writer.checkpoint().expect("save interrupted enumeration"); + let encoded = rmp_serde::to_vec(&saved).expect("encode checkpoint"); + let restored = rmp_serde::from_slice(&encoded).expect("restore checkpoint"); + compare_with_cumulative_ingestion(restored, &["h", "g", "f", "e", "d", "c", "b", "a", "i", "j"]); + } + } +} + +#[test] +fn raw_enumeration_writer_preserves_restored_building_and_terminal_pages() { + for complete in [false, true] { + for budget in [1, 2] { + let mut index = RawEnumerationPageIndex::new("bucket", 2).expect("index"); + let source = ["a".to_owned(), "b".to_owned()]; + if complete { + index.ingest_owner_entries(source, budget, 0).expect("complete source"); + } else { + index.ingest_partial_owner_entries(source, budget, 0).expect("partial source"); + } + compare_with_cumulative_ingestion(index.clone(), &["b", "a", "a", "c", "d"]); + compare_with_cumulative_ingestion(checkpoint(index), &["b", "a", "a", "c", "d"]); + } + } +} + +#[test] +fn raw_enumeration_writer_keeps_identity_check_when_terminal_builder_commits() { + for limit in [2, 3] { + let mut index = RawEnumerationPageIndex::new("bucket", limit).expect("index"); + index + .ingest_owner_entries(["a".to_owned(), "b".to_owned()], 2, 0) + .expect("terminal building page"); + compare_with_cumulative_ingestion(index.clone(), &["b", "a", "b", "c"]); + compare_with_cumulative_ingestion(index, &["foreign", "a", "b", "c"]); + } +} + +#[test] +fn raw_enumeration_writer_rejects_complete_source_drift_after_observation_floor() { + let mut index = RawEnumerationPageIndex::new("bucket", 2).expect("index"); + index + .ingest_owner_entries(["a".to_owned(), "b".to_owned()], 2, 0) + .expect("complete source"); + let index = checkpoint(index); + for source in [["foreign", "a"], ["a", "foreign"], ["a", "a"], ["b", "a"]] { + compare_with_cumulative_ingestion(index.clone(), &source); + } +} + +#[test] +fn raw_enumeration_writer_rejects_corruption_at_restore() { + let mut writer = RawEnumerationPageWriter::new(RawEnumerationPageIndex::new("bucket", 2).expect("index")).expect("writer"); + for entry in ["a", "b", "c"] { + writer.record_entry(entry).expect("valid entry"); + } + let valid = writer.checkpoint().expect("checkpoint"); + for case in 0..7 { + let mut corrupt = valid.clone(); + let RawEnumerationPageIndexState::Supported(inner) = &mut corrupt.state else { + panic!("supported checkpoint"); + }; + match case { + 0 => inner.pages[0].digest[0] ^= 1, + 1 => inner.pages[1].entries_start = 0, + 2 => inner.pages[1] = inner.pages[0].clone(), + 3 => inner.page_entry_limit = 0, + 4 => inner.complete = true, + 5 => inner.parent.clear(), + 6 => { + inner.building = Some(RawEnumerationPageBuilder { + page_index: 2, + entries_start: 3, + entries: vec!["a".to_owned()], + terminal: false, + }); + } + _ => unreachable!(), + } + let bytes = rmp_serde::to_vec(&corrupt).expect("encode damaged checkpoint"); + let restored = rmp_serde::from_slice(&bytes).expect("decode untrusted checkpoint"); + assert!( + matches!(RawEnumerationPageWriter::new(restored), Err(RawEnumerationPageIndexError::CorruptIndex)), + "corruption case {case} must not become trusted runtime state" + ); + } +} + +#[test] +fn raw_enumeration_writer_rejects_invalid_names_without_publishing_them() { + for entry in ["", ".", "..", "nested/name", &"x".repeat(16 * 1024 + 1)] { + let mut writer = + RawEnumerationPageWriter::new(RawEnumerationPageIndex::new("bucket", 128).expect("index")).expect("writer"); + let before = writer.checkpoint().expect("empty checkpoint"); + assert_eq!(writer.record_entry(entry), Err(RawEnumerationPageIndexError::InvalidEntry)); + assert_eq!(writer.checkpoint().expect("checkpoint after rejection"), before); + assert_eq!(writer.indexed_entry_count(), 0); + } +} + +#[test] +fn raw_enumeration_writer_hashes_only_new_pages_after_restore() { + let count = 4096; + let mut writer = RawEnumerationPageWriter::new(RawEnumerationPageIndex::new("bucket", 128).expect("index")).expect("writer"); + RAW_PAGE_DIGEST_ENTRIES.set(0); + for i in (0..count).rev() { + writer.record_entry(&format!("entry-{i:06}")).expect("append entry"); + } + assert_eq!( + RAW_PAGE_DIGEST_ENTRIES.get(), + count, + "each committed entry is hashed once during enumeration" + ); + let saved = writer.checkpoint().expect("checkpoint"); + assert_eq!(saved.indexed_entries().expect("validated checkpoint").len(), count); + let mut resumed = RawEnumerationPageWriter::new(saved).expect("restore large checkpoint"); + RAW_PAGE_DIGEST_ENTRIES.set(0); + for i in 0..count * 2 { + resumed.record_entry(&format!("entry-{i:06}")).expect("resume and append"); + } + assert_eq!(resumed.indexed_entry_count(), count * 2); + assert_eq!( + RAW_PAGE_DIGEST_ENTRIES.get(), + count, + "restored pages must not be rehashed for each entry or commit" + ); +} + +#[test] +fn raw_enumeration_writer_keeps_snapshot_independent_of_later_appends() { + let mut writer = RawEnumerationPageWriter::new(RawEnumerationPageIndex::new("bucket", 128).expect("index")).expect("writer"); + for i in 0..127 { + writer.record_entry(&format!("entry-{i:03}")).expect("append before boundary"); + } + let partial = writer.checkpoint().expect("partial page checkpoint"); + writer.record_entry("entry-127").expect("full page"); + let full = writer.checkpoint().expect("full page checkpoint"); + writer.record_entry("entry-128").expect("next page"); + assert_eq!(partial.indexed_entries().expect("old partial snapshot").len(), 127); + assert_eq!(full.indexed_entries().expect("old full snapshot").len(), 128); + assert_eq!( + writer + .checkpoint() + .expect("new snapshot") + .indexed_entries() + .expect("new entries") + .len(), + 129 + ); + assert_ne!(partial, full); +} diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index ccb748f11..ff1d80d4d 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -25,7 +25,7 @@ use crate::data_usage_define::{ PendingScannerHealKind, ScannerSizeSummaryExt, SizeReconciliationEntry, SizeSummary, hash_path, }; use crate::error::ScannerError; -use crate::raw_page_index::{RawEnumerationPageIndex, RawEnumerationPageIndexError}; +use crate::raw_page_index::{RawEnumerationPageIndex, RawEnumerationPageWriter}; use crate::runtime_config::{ scanner_alert_excess_folders, scanner_alert_excess_version_size, scanner_alert_excess_versions, scanner_yield_every_n_objects, }; @@ -95,7 +95,6 @@ const SCANNER_CHECKPOINT_OBJECT_INTERVAL: u64 = 1024; const SCANNER_CHECKPOINT_MIN_INTERVAL: Duration = Duration::from_secs(5); const SCANNER_CHECKPOINT_INTERVAL: Duration = Duration::from_secs(60); const SCANNER_RAW_ENUMERATION_PAGE_ENTRY_LIMIT: usize = 128; -const SCANNER_RAW_ENUMERATION_PAGE_BUILD_BUDGET: usize = 1; // Erasure data directories contain direct part.N files; keep namespace probes bounded. const ERASURE_DATA_DIR_PROBE_ENTRY_LIMIT: usize = 64; const DEFAULT_HEAL_OBJECT_SELECT_PROB: u32 = 1024; @@ -761,33 +760,21 @@ struct RawEnumerationProgress { last_entry: Option, entries_seen: u64, digest: Sha256, - observed_entries: Vec, - revalidate_after_entries: usize, - page_index: Option, + page_index: Option, } impl RawEnumerationProgress { fn new(parent: &str, page_index: Option) -> Self { let mut digest = Sha256::new(); update_raw_enumeration_digest(&mut digest, b"parent", parent.as_bytes()); - let mut revalidate_after_entries = 0; - let page_index = match page_index { - Some(index) => match index.indexed_entries() { - Ok(entries) => { - revalidate_after_entries = entries.len(); - Some(index) - } - Err(_) => None, - }, - None => RawEnumerationPageIndex::new(parent, SCANNER_RAW_ENUMERATION_PAGE_ENTRY_LIMIT).ok(), - }; + let page_index = page_index + .or_else(|| RawEnumerationPageIndex::new(parent, SCANNER_RAW_ENUMERATION_PAGE_ENTRY_LIMIT).ok()) + .and_then(|index| RawEnumerationPageWriter::new(index).ok()); Self { parent: parent.to_string(), last_entry: None, entries_seen: 0, digest, - observed_entries: Vec::new(), - revalidate_after_entries, page_index, } } @@ -796,34 +783,10 @@ impl RawEnumerationProgress { update_raw_enumeration_digest(&mut self.digest, b"entry", entry.as_bytes()); self.last_entry = Some(entry.to_string()); self.entries_seen = self.entries_seen.saturating_add(1); - self.observed_entries.push(entry.to_string()); - if let Some(index) = &mut self.page_index { - if self.observed_entries.len() < self.revalidate_after_entries { - return; - } - let result = index - .generation() - .ok_or(RawEnumerationPageIndexError::Unsupported) - .and_then(|generation| { - index.ingest_partial_owner_entries( - self.observed_entries.clone(), - SCANNER_RAW_ENUMERATION_PAGE_BUILD_BUDGET, - generation, - ) - }); - match result { - Ok(outcome) if outcome.ready_to_commit => { - if let Some(generation) = index.generation() - && index.commit_building_page(generation).is_err() - { - self.page_index = None; - } - } - Ok(_) => {} - Err(_) => { - self.page_index = None; - } - } + if let Some(index) = &mut self.page_index + && index.record_entry(entry).is_err() + { + self.page_index = None; } } @@ -840,28 +803,20 @@ impl RawEnumerationProgress { } fn page_index(&self) -> Option { - self.page_index.clone().and_then(|mut index| { - if let Some(generation) = index.generation() - && matches!(index.status(), crate::raw_page_index::RawEnumerationPageOwnerStatus::Building { .. }) - && index.commit_building_page(generation).is_err() - { - return None; - } - match index.indexed_entries() { - Ok(entries) if !entries.is_empty() => Some(index), - _ => None, - } - }) + self.page_index + .as_ref() + .filter(|index| index.indexed_entry_count() > 0) + .and_then(|index| index.checkpoint().ok()) } fn has_checkpointable_page_index(&self) -> bool { - self.page_index().is_some() + self.checkpointable_entry_count() > 0 } fn checkpointable_entry_count(&self) -> usize { - self.page_index() - .and_then(|index| index.indexed_entries().ok()) - .map_or(0, |entries| entries.len()) + self.page_index + .as_ref() + .map_or(0, RawEnumerationPageWriter::indexed_entry_count) } } @@ -1130,19 +1085,6 @@ impl FolderScanner { if self.old_cache.info.scan_progress.is_none() { return; } - let page_index = self - .old_cache - .validated_raw_enumeration_page_index() - .filter(|index| match index.status() { - crate::raw_page_index::RawEnumerationPageOwnerStatus::Building { - parent: index_parent, .. - } - | crate::raw_page_index::RawEnumerationPageOwnerStatus::Ready { - parent: index_parent, .. - } => index_parent == parent, - crate::raw_page_index::RawEnumerationPageOwnerStatus::Unsupported => false, - }) - .cloned(); if let Some(position) = self .raw_enumeration_progress .iter() @@ -1150,6 +1092,7 @@ impl FolderScanner { { self.raw_enumeration_progress.truncate(position + 1); } else { + let page_index = self.raw_enumeration_page_index_for_parent(parent).cloned(); self.raw_enumeration_progress .push(RawEnumerationProgress::new(parent, page_index)); } @@ -1158,8 +1101,17 @@ impl FolderScanner { } } + fn raw_enumeration_page_index_for_parent(&self, parent: &str) -> Option<&RawEnumerationPageIndex> { + self.old_cache + .info + .scan_raw_enumeration_page_index + .as_ref() + .filter(|index| index.parent() == Some(parent))?; + self.old_cache.validated_raw_enumeration_page_index() + } + fn raw_enumeration_committed_entry_oracle(&self, parent: &str) -> HashSet { - let Some(index) = self.old_cache.validated_raw_enumeration_page_index() else { + let Some(index) = self.raw_enumeration_page_index_for_parent(parent) else { return HashSet::new(); }; let generation_matches_parent = match index.status() { diff --git a/crates/scanner/src/scanner_folder/tests.rs b/crates/scanner/src/scanner_folder/tests.rs index a940a2eb2..fea7f6499 100644 --- a/crates/scanner/src/scanner_folder/tests.rs +++ b/crates/scanner/src/scanner_folder/tests.rs @@ -28,6 +28,7 @@ use std::sync::Mutex; mod checkpoint_fixture; pub(super) mod enumeration_restart; +mod incremental_enumeration; /// Reset the process-global alert cooldown map; test-only. fn reset_alert_cooldowns() { diff --git a/crates/scanner/src/scanner_folder/tests/incremental_enumeration.rs b/crates/scanner/src/scanner_folder/tests/incremental_enumeration.rs new file mode 100644 index 000000000..9390306db --- /dev/null +++ b/crates/scanner/src/scanner_folder/tests/incremental_enumeration.rs @@ -0,0 +1,93 @@ +// Copyright 2024 RustFS Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use super::*; +use crate::raw_page_index::RAW_PAGE_DIGEST_ENTRIES; + +#[tokio::test] +async fn raw_enumeration_reuses_validated_parent_and_skips_unrelated_index() { + let (mut scanner, temp_dir) = build_test_scanner().await; + let _guard = TestGuard { + temp_dir: Some(temp_dir), + }; + scanner.old_cache.info.name = "bucket".to_owned(); + scanner.old_cache.info.scan_progress = Some(crate::DataUsageScanProgress { + started_plan: crate::DataUsageScanPlanDigest([1; 32]), + requested_plan: crate::DataUsageScanPlanDigest([1; 32]), + }); + scanner.old_cache.info.source = Some(crate::DataUsageCacheSource::new(0, 0)); + scanner.old_cache.info.scan_identity = Some(crate::DataUsageScanIdentity { + version: 1, + bucket_incarnation: Uuid::from_u128(1), + set_layout: crate::DataUsageScanPlanDigest([2; 32]), + publication_epoch: 1, + tier_registry_generation: 0, + scan_mode: HealScanMode::Normal, + }); + let mut saved = + RawEnumerationPageWriter::new(RawEnumerationPageIndex::new("bucket/metadata", 128).expect("index")).expect("writer"); + for i in 0..1024 { + saved.record_entry(&format!("entry-{i:04}")).expect("initial entry"); + } + scanner.old_cache.info.scan_raw_enumeration_page_index = Some(saved.checkpoint().expect("saved index")); + scanner.record_raw_enumeration_entry("bucket/metadata", "entry-0000"); + RAW_PAGE_DIGEST_ENTRIES.set(0); + for i in 1..1024 { + scanner.record_raw_enumeration_entry("bucket/metadata", &format!("entry-{i:04}")); + } + assert_eq!( + RAW_PAGE_DIGEST_ENTRIES.get(), + 0, + "the same parent must not revalidate its restored index per entry" + ); + let progress = scanner.raw_enumeration_progress.last().expect("active progress"); + assert_eq!(progress.checkpointable_entry_count(), 1024); + assert!(progress.has_checkpointable_page_index()); + assert_eq!( + RAW_PAGE_DIGEST_ENTRIES.get(), + 0, + "progress selection must use cached counts, not snapshot/validation" + ); + for i in 0..32 { + let parent = format!("bucket/other-{i}"); + assert!(scanner.raw_enumeration_committed_entry_oracle(&parent).is_empty()); + scanner.record_raw_enumeration_entry(&parent, "object"); + scanner.finish_raw_enumeration_parent(&parent); + } + assert_eq!( + RAW_PAGE_DIGEST_ENTRIES.get(), + 0, + "unrelated directories must not validate the saved metadata index" + ); + let (_, checkpoint) = scanner.take_raw_enumeration_resume_state(); + assert_eq!( + checkpoint + .expect("metadata checkpoint") + .indexed_entries() + .expect("valid restored output") + .len(), + 1024 + ); +} + +#[test] +fn raw_enumeration_progress_invalid_entry_discards_index_but_retains_diagnostic_cursor() { + let mut progress = RawEnumerationProgress::new("bucket", None); + progress.record_entry("valid"); + progress.record_entry("invalid/name"); + progress.record_entry("later"); + assert_eq!(progress.checkpointable_entry_count(), 0); + assert!(progress.page_index().is_none()); + assert_eq!(progress.cursor().expect("diagnostic cursor").entries_seen, 3); +}