From fe3e7799778a787d7f227f4c430b6b7c847f689c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Sun, 23 Aug 2026 05:45:30 +0800 Subject: [PATCH] fix(scanner): fence unknown tier accounting --- crates/data-usage/src/data_usage.rs | 533 ++++++++++++++++-- crates/ecstore/src/api/mod.rs | 12 +- crates/ecstore/src/cluster/rpc/http_auth.rs | 81 ++- .../cluster/rpc/internode_data_transport.rs | 229 +++++++- crates/ecstore/src/cluster/rpc/mod.rs | 5 +- crates/ecstore/src/cluster/rpc/remote_disk.rs | 61 +- .../ecstore/src/storage_api_contracts/mod.rs | 12 +- crates/scanner/src/data_usage_define.rs | 126 ++++- crates/scanner/src/data_usage_define/tests.rs | 100 +++- crates/scanner/src/lib.rs | 241 +++++++- crates/scanner/src/remote_scanner/stream.rs | 39 +- .../src/remote_scanner/stream/tests.rs | 38 ++ crates/scanner/src/scanner_folder.rs | 19 +- .../src/scanner_folder/item_actions.rs | 73 ++- crates/scanner/src/scanner_io.rs | 21 +- crates/scanner/src/scanner_io/cache.rs | 186 +++++- crates/scanner/src/scanner_io/io_cache.rs | 24 +- crates/scanner/src/scanner_io/io_cycle.rs | 11 + crates/scanner/src/scanner_io/io_disk.rs | 11 +- .../src/scanner_io/publish_gate_tests.rs | 340 ++++++++++- crates/scanner/src/scanner_io/tests.rs | 43 +- crates/storage-api/src/lib.rs | 3 + rustfs/src/storage/rpc/http_service.rs | 113 +++- rustfs/src/storage/storage_api.rs | 27 +- 24 files changed, 2128 insertions(+), 220 deletions(-) diff --git a/crates/data-usage/src/data_usage.rs b/crates/data-usage/src/data_usage.rs index ca00edc32..e7542933a 100644 --- a/crates/data-usage/src/data_usage.rs +++ b/crates/data-usage/src/data_usage.rs @@ -12,8 +12,13 @@ // See the License for the specific language governing permissions and // limitations under the License. -use serde::{Deserialize, Serialize, ser::SerializeMap as _}; +use serde::{ + Deserialize, Serialize, + de::{IgnoredAny, SeqAccess, Visitor}, + ser::SerializeMap as _, +}; use std::{ + borrow::Cow, collections::{HashMap, HashSet}, hash::{DefaultHasher, Hash, Hasher}, time::{Duration, SystemTime}, @@ -105,8 +110,9 @@ pub const UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP: usize = 64; pub const UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP: usize = 4096; pub const UNKNOWN_TIER_DIAGNOSTIC_TTL: Duration = Duration::from_secs(60 * 60); const UNKNOWN_TIER_DIAGNOSTIC_KEY_BYTES: usize = 256; +const UNKNOWN_TIER_DIAGNOSTIC_INPUT_CAP: usize = UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP * 4; -#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)] +#[derive(Clone, Debug, Default, Serialize, PartialEq, Eq)] pub struct UnknownTierStats { /// Logical bytes retained in the scanner's normal usage total. pub unknown_bytes: u64, @@ -123,6 +129,140 @@ pub struct UnknownTierStats { pub diagnostics: Vec, #[serde(default)] pub diagnostics_at: Option, + /// A saturating update occurred; this snapshot cannot prove conservation. + /// The field is persisted so a restarted scanner cannot mistake a + /// saturated legacy aggregate for exact accounting evidence. + #[serde(default)] + pub counter_overflowed: bool, +} + +#[derive(Default)] +struct BoundedDiagnostics { + entries: Vec, + dropped: u64, +} + +impl<'de> Deserialize<'de> for BoundedDiagnostics { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + struct BoundedDiagnosticsVisitor; + + impl<'de> Visitor<'de> for BoundedDiagnosticsVisitor { + type Value = BoundedDiagnostics; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("a sequence of bounded tier diagnostics") + } + + fn visit_seq(self, mut sequence: A) -> Result + where + A: SeqAccess<'de>, + { + let mut bounded = BoundedDiagnostics::default(); + let mut bytes = 0_usize; + let mut inspected = 0_usize; + loop { + if bounded.entries.len() >= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP + || bytes >= UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP + || inspected >= UNKNOWN_TIER_DIAGNOSTIC_INPUT_CAP + { + if sequence.next_element::()?.is_none() { + break; + } + bounded.dropped = bounded.dropped.saturating_add(1); + continue; + } + let Some(diagnostic) = sequence.next_element::>()? else { + break; + }; + inspected = inspected.saturating_add(1); + if !is_redacted_tier_diagnostic(&diagnostic) + || bytes.saturating_add(diagnostic.len()) > UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP + || bounded.entries.iter().any(|entry| entry == &diagnostic) + { + bounded.dropped = bounded.dropped.saturating_add(1); + continue; + } + bytes = bytes.saturating_add(diagnostic.len()); + bounded.entries.push(diagnostic.into_owned()); + } + Ok(bounded) + } + } + + deserializer.deserialize_seq(BoundedDiagnosticsVisitor) + } +} + +#[derive(Deserialize)] +struct UnknownTierStatsWire { + #[serde(default)] + unknown_bytes: u64, + #[serde(default)] + unknown_physical_bytes: u64, + #[serde(default)] + unknown_objects: u64, + #[serde(default)] + unknown_versions: u64, + #[serde(default)] + diagnostics_dropped: u64, + #[serde(default)] + diagnostics: BoundedDiagnostics, + #[serde(default)] + diagnostics_at: Option, + #[serde(default)] + counter_overflowed: bool, +} + +fn is_redacted_tier_diagnostic(diagnostic: &str) -> bool { + diagnostic + .strip_prefix("tier-hash:") + .is_some_and(|digest| digest.len() == 16 && digest.bytes().all(|byte| byte.is_ascii_hexdigit())) +} + +fn diagnostics_expired(at: SystemTime, now: SystemTime) -> bool { + now.duration_since(at).map_or(true, |age| age > UNKNOWN_TIER_DIAGNOSTIC_TTL) +} + +fn checked_saturating_add(left: u64, right: u64, overflowed: &mut bool) -> u64 { + match left.checked_add(right) { + Some(value) => value, + None => { + *overflowed = true; + u64::MAX + } + } +} + +impl<'de> Deserialize<'de> for UnknownTierStats { + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + let wire = UnknownTierStatsWire::deserialize(deserializer)?; + let mut stats = Self { + unknown_bytes: wire.unknown_bytes, + unknown_physical_bytes: wire.unknown_physical_bytes, + unknown_objects: wire.unknown_objects, + unknown_versions: wire.unknown_versions, + counter_overflowed: wire.counter_overflowed, + diagnostics_dropped: wire.diagnostics_dropped, + diagnostics: wire.diagnostics.entries, + diagnostics_at: wire.diagnostics_at, + }; + stats.diagnostics_dropped = + checked_saturating_add(stats.diagnostics_dropped, wire.diagnostics.dropped, &mut stats.counter_overflowed); + if stats + .diagnostics_at + .is_some_and(|at| diagnostics_expired(at, SystemTime::now())) + { + stats.diagnostics.clear(); + stats.diagnostics_at = None; + } + Ok(stats) + } } impl UnknownTierStats { @@ -137,25 +277,24 @@ impl UnknownTierStats { /// tier bytes. Both counters are saturating so malformed metadata cannot /// wrap an aggregate. pub fn record_dimensions(&mut self, tier: &str, logical_bytes: u64, physical_bytes: u64, versions: u64, objects: u64) { - self.unknown_bytes = self.unknown_bytes.saturating_add(logical_bytes); - self.unknown_physical_bytes = self.unknown_physical_bytes.saturating_add(physical_bytes); - self.unknown_objects = self.unknown_objects.saturating_add(objects); - self.unknown_versions = self.unknown_versions.saturating_add(versions); + self.unknown_bytes = checked_saturating_add(self.unknown_bytes, logical_bytes, &mut self.counter_overflowed); + self.unknown_physical_bytes = + checked_saturating_add(self.unknown_physical_bytes, physical_bytes, &mut self.counter_overflowed); + self.unknown_objects = checked_saturating_add(self.unknown_objects, objects, &mut self.counter_overflowed); + self.unknown_versions = checked_saturating_add(self.unknown_versions, versions, &mut self.counter_overflowed); let digest = { let mut hasher = DefaultHasher::new(); // Bound hashing work for hostile metadata while retaining enough // length/prefix entropy to reconcile repeated observations. - hasher.write_usize(tier.len()); - hasher.write(&tier.as_bytes()[..tier.len().min(UNKNOWN_TIER_DIAGNOSTIC_KEY_BYTES)]); - hasher.write_u8(u8::from(tier.bytes().any(|byte| byte.is_ascii_control()))); + hasher.write_u64(u64::try_from(tier.len()).unwrap_or(u64::MAX)); + let bounded = &tier.as_bytes()[..tier.len().min(UNKNOWN_TIER_DIAGNOSTIC_KEY_BYTES)]; + hasher.write(bounded); + hasher.write_u8(u8::from(bounded.iter().any(|byte| byte.is_ascii_control()))); format!("tier-hash:{:016x}", hasher.finish()) }; let now = SystemTime::now(); - if self - .diagnostics_at - .is_some_and(|at| now.duration_since(at).unwrap_or_default() > UNKNOWN_TIER_DIAGNOSTIC_TTL) - { + if self.diagnostics_at.is_some_and(|at| diagnostics_expired(at, now)) { self.diagnostics.clear(); } self.diagnostics_at = Some(now); @@ -166,44 +305,55 @@ impl UnknownTierStats { if self.diagnostics.len() >= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP || current_bytes.saturating_add(digest.len()) > UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP { - self.diagnostics_dropped = self.diagnostics_dropped.saturating_add(1); + self.diagnostics_dropped = checked_saturating_add(1, self.diagnostics_dropped, &mut self.counter_overflowed); return; } self.diagnostics.push(digest); } pub fn merge(&mut self, other: &Self) { - self.unknown_bytes = self.unknown_bytes.saturating_add(other.unknown_bytes); - self.unknown_physical_bytes = self.unknown_physical_bytes.saturating_add(other.unknown_physical_bytes); - self.unknown_objects = self.unknown_objects.saturating_add(other.unknown_objects); - self.unknown_versions = self.unknown_versions.saturating_add(other.unknown_versions); - self.diagnostics_dropped = self.diagnostics_dropped.saturating_add(other.diagnostics_dropped); + self.unknown_bytes = checked_saturating_add(self.unknown_bytes, other.unknown_bytes, &mut self.counter_overflowed); + self.unknown_physical_bytes = + checked_saturating_add(self.unknown_physical_bytes, other.unknown_physical_bytes, &mut self.counter_overflowed); + self.unknown_objects = checked_saturating_add(self.unknown_objects, other.unknown_objects, &mut self.counter_overflowed); + self.unknown_versions = + checked_saturating_add(self.unknown_versions, other.unknown_versions, &mut self.counter_overflowed); + self.diagnostics_dropped = + checked_saturating_add(self.diagnostics_dropped, other.diagnostics_dropped, &mut self.counter_overflowed); + self.counter_overflowed |= other.counter_overflowed; let now = SystemTime::now(); - if self - .diagnostics_at - .zip(other.diagnostics_at) - .is_some_and(|(left, right)| now.duration_since(left.max(right)).unwrap_or_default() > UNKNOWN_TIER_DIAGNOSTIC_TTL) - { + if self.diagnostics_at.is_some_and(|at| diagnostics_expired(at, now)) { self.diagnostics.clear(); + self.diagnostics_at = None; } - self.diagnostics_at = Some(now); - for diagnostic in &other.diagnostics { - if self.diagnostics.iter().any(|entry| entry == diagnostic) { - continue; + if !other.diagnostics_at.is_some_and(|at| diagnostics_expired(at, now)) { + for diagnostic in &other.diagnostics { + if !is_redacted_tier_diagnostic(diagnostic) { + self.diagnostics_dropped = checked_saturating_add(1, self.diagnostics_dropped, &mut self.counter_overflowed); + continue; + } + if self.diagnostics.iter().any(|entry| entry == diagnostic) { + continue; + } + let current_bytes: usize = self.diagnostics.iter().map(String::len).sum(); + if self.diagnostics.len() >= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP + || current_bytes.saturating_add(diagnostic.len()) > UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP + { + self.diagnostics_dropped = checked_saturating_add(1, self.diagnostics_dropped, &mut self.counter_overflowed); + continue; + } + self.diagnostics.push(diagnostic.clone()); } - let current_bytes: usize = self.diagnostics.iter().map(String::len).sum(); - if self.diagnostics.len() >= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP - || current_bytes.saturating_add(diagnostic.len()) > UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP - { - self.diagnostics_dropped = self.diagnostics_dropped.saturating_add(1); - continue; - } - self.diagnostics.push(diagnostic.clone()); + } + if !self.diagnostics.is_empty() { + self.diagnostics_at = Some(now); } } pub fn fits_add(&self, other: &Self) -> bool { - self.unknown_bytes.checked_add(other.unknown_bytes).is_some() + !self.counter_overflowed + && !other.counter_overflowed + && self.unknown_bytes.checked_add(other.unknown_bytes).is_some() && self .unknown_physical_bytes .checked_add(other.unknown_physical_bytes) @@ -213,8 +363,18 @@ impl UnknownTierStats { && self.diagnostics_dropped.checked_add(other.diagnostics_dropped).is_some() } + pub fn checked_add(&self, other: &Self) -> Option { + if !self.fits_add(other) { + return None; + } + let mut merged = self.clone(); + merged.merge(other); + Some(merged) + } + pub fn is_empty(&self) -> bool { - self.unknown_bytes == 0 + !self.counter_overflowed + && self.unknown_bytes == 0 && self.unknown_physical_bytes == 0 && self.unknown_objects == 0 && self.unknown_versions == 0 @@ -273,13 +433,13 @@ impl AllTierStats { { let known: HashSet<&str> = known_tiers.into_iter().collect(); let mut unknown = self.tiers.remove(UNKNOWN_TIER).unwrap_or_default(); - let retired: Vec = self + let retired_tiers: Vec = self .tiers .keys() .filter(|tier| tier.as_str() != "STANDARD" && tier.as_str() != "REDUCED_REDUNDANCY" && !known.contains(tier.as_str())) .cloned() .collect(); - for tier in retired { + for tier in retired_tiers { if let Some(stats) = self.tiers.remove(&tier) { unknown = unknown.add(&stats); } @@ -475,6 +635,48 @@ pub struct DiskUsageStatus { pub snapshot_exists: bool, } +/// Independent conservation evidence for a scanner summary. +/// +/// The totals are maintained while objects are accounted and are deliberately +/// not reconstructed from the tier map at publish time. `*_known` records the +/// portion represented by the corresponding accounting dimension. +#[derive(Debug, Default, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] +pub struct TierAccountingProof { + pub logical_total: u64, + pub logical_known: u64, + pub physical_total: u64, + pub physical_known: u64, + #[serde(default)] + pub overflowed: bool, +} + +impl TierAccountingProof { + pub fn checked_add(self, other: Self) -> Option { + if self.overflowed || other.overflowed { + return None; + } + Some(Self { + logical_total: self.logical_total.checked_add(other.logical_total)?, + logical_known: self.logical_known.checked_add(other.logical_known)?, + physical_total: self.physical_total.checked_add(other.physical_total)?, + physical_known: self.physical_known.checked_add(other.physical_known)?, + overflowed: false, + }) + } + + pub fn saturating_add(&mut self, other: Self) { + let Some(merged) = (*self).checked_add(other) else { + self.logical_total = self.logical_total.saturating_add(other.logical_total); + self.logical_known = self.logical_known.saturating_add(other.logical_known); + self.physical_total = self.physical_total.saturating_add(other.physical_total); + self.physical_known = self.physical_known.saturating_add(other.physical_known); + self.overflowed = true; + return; + }; + *self = merged; + } +} + /// Size summary for a single object or group of objects #[derive(Debug, Default, Clone)] pub struct SizeSummary { @@ -506,6 +708,8 @@ pub struct SizeSummary { pub tier_stats: HashMap, /// Counters and bounded diagnostics for unknown tiers in this summary. pub unknown_tier_stats: UnknownTierStats, + /// Independent logical/physical conservation evidence. + pub tier_accounting_proof: TierAccountingProof, } /// Replication target size summary @@ -597,6 +801,10 @@ impl<'de> Deserialize<'de> for SizeHistogram { } impl SizeHistogram { + pub fn is_empty(&self) -> bool { + self.0.iter().all(|value| *value == 0) + } + pub fn add(&mut self, size: u64) { let intervals = [ (0, 1024 - 1), // LESS_THAN_1024_B @@ -711,6 +919,10 @@ impl<'de> Deserialize<'de> for VersionsHistogram { } impl VersionsHistogram { + pub fn is_empty(&self) -> bool { + self.0.iter().all(|value| *value == 0) + } + pub fn add(&mut self, count: u64) { let intervals = [ (0, 0), // UNVERSIONED @@ -854,6 +1066,10 @@ pub struct DataUsageEntry { /// Bounded unknown-tier reconciliation state for this cache entry. #[serde(default)] pub unknown_tier_stats: Option, + /// Optional conservation proof. Missing values are legacy/unproven, not + /// zero-valued evidence. + #[serde(default)] + pub tier_accounting_proof: Option, } impl Serialize for DataUsageEntry { @@ -864,7 +1080,9 @@ impl Serialize for DataUsageEntry { // Keep entries map-encoded so older readers can ignore fields appended // by newer scanner versions during rolling upgrades. The derived // (array) encoding made any appended field a decode error for them. - let mut state = serializer.serialize_map(Some(12))?; + let mut state = serializer.serialize_map(Some( + 11 + usize::from(self.unknown_tier_stats.is_some()) + usize::from(self.tier_accounting_proof.is_some()), + ))?; state.serialize_entry("children", &self.children)?; state.serialize_entry("size", &self.size)?; state.serialize_entry("objects", &self.objects)?; @@ -876,12 +1094,35 @@ impl Serialize for DataUsageEntry { state.serialize_entry("compacted", &self.compacted)?; state.serialize_entry("failed_objects", &self.failed_objects)?; state.serialize_entry("all_tier_stats", &self.all_tier_stats)?; - state.serialize_entry("unknown_tier_stats", &self.unknown_tier_stats)?; + // Keep the legacy no-unknown shape byte-for-byte stable. Once unknown + // accounting exists, append its map field before the optional proof. + if let Some(unknown_tier_stats) = self.unknown_tier_stats.as_ref() { + state.serialize_entry("unknown_tier_stats", unknown_tier_stats)?; + } + if let Some(proof) = self.tier_accounting_proof.as_ref() { + state.serialize_entry("tier_accounting_proof", proof)?; + } state.end() } } impl DataUsageEntry { + fn has_local_usage(&self) -> bool { + self.size != 0 + || self.objects != 0 + || self.versions != 0 + || self.delete_markers != 0 + || self.failed_objects != 0 + || self.obj_sizes.0.iter().any(|value| *value != 0) + || self.obj_versions.0.iter().any(|value| *value != 0) + || self.replication_stats.as_ref().is_some_and(|stats| !stats.is_empty()) + || self + .all_tier_stats + .as_ref() + .is_some_and(|stats| stats.tiers.values().any(|tier| !tier.is_empty())) + || self.unknown_tier_stats.as_ref().is_some_and(|stats| !stats.is_empty()) + } + pub fn add_child(&mut self, hash: &DataUsageHash) { if self.children.contains(&hash.key()) { return; @@ -890,6 +1131,8 @@ impl DataUsageEntry { } pub fn merge(&mut self, other: &DataUsageEntry) { + let self_had_local_usage = self.has_local_usage(); + let other_had_local_usage = other.has_local_usage(); self.objects += other.objects; self.versions += other.versions; self.delete_markers += other.delete_markers; @@ -924,6 +1167,24 @@ impl DataUsageEntry { .merge(other_unknown); } + self.tier_accounting_proof = match (self.tier_accounting_proof, other.tier_accounting_proof) { + (Some(mut left), Some(right)) => { + left.saturating_add(right); + Some(left) + } + (None, Some(right)) if !self_had_local_usage => Some(right), + (Some(left), None) if !other_had_local_usage => Some(left), + (None, None) => None, + _ => None, + }; + // A saturated unknown-tier counter invalidates conservation evidence; + // never let an otherwise valid proof hide that loss of precision. + if self.unknown_tier_stats.as_ref().is_some_and(|stats| stats.counter_overflowed) + && let Some(proof) = self.tier_accounting_proof.as_mut() + { + proof.overflowed = true; + } + self.obj_sizes.merge_from(&other.obj_sizes); self.obj_versions.merge_from(&other.obj_versions); } @@ -1009,11 +1270,20 @@ impl DataUsageEntry { (Some(left), Some(right)) => left.fits_merge(right), }; let unknown_tier_stats_fit = match (&self.unknown_tier_stats, &other.unknown_tier_stats) { - (_, None) | (None, Some(_)) => true, + (None, None) => true, + (Some(left), None) => !left.counter_overflowed, + (None, Some(right)) => !right.counter_overflowed, (Some(left), Some(right)) => left.fits_add(right), }; - if !scalar_counts_fit || !histograms_fit || !replication_fits || !tier_stats_fit || !unknown_tier_stats_fit { + let proof_fit = match (self.tier_accounting_proof, other.tier_accounting_proof) { + (Some(left), Some(right)) => left.checked_add(right).is_some(), + (Some(proof), None) | (None, Some(proof)) => !proof.overflowed, + (None, None) => true, + }; + + if !scalar_counts_fit || !histograms_fit || !replication_fits || !tier_stats_fit || !unknown_tier_stats_fit || !proof_fit + { return false; } self.merge(other); @@ -1531,7 +1801,7 @@ impl DataUsageCache { delete_markers_total_count: flat.delete_markers as u64, objects_total_size: flat.size as u64, tier_stats: flat.all_tier_stats.filter(|tiers| !tiers.is_empty()), - unknown_tier_stats: flat.unknown_tier_stats, + unknown_tier_stats: flat.unknown_tier_stats.filter(|stats| !stats.is_empty()), buckets_count: u64::try_from(buckets.len()).unwrap_or(u64::MAX), buckets_usage, usage_snapshot_complete: self.info.snapshot_complete, @@ -1960,6 +2230,10 @@ impl SizeSummary { self.pending_count = self.pending_count.saturating_add(other.pending_count); self.failed_count = self.failed_count.saturating_add(other.failed_count); self.unknown_tier_stats.merge(&other.unknown_tier_stats); + self.tier_accounting_proof.saturating_add(other.tier_accounting_proof); + if self.unknown_tier_stats.counter_overflowed { + self.tier_accounting_proof.overflowed = true; + } // A disk/bucket aggregate is assembled from many object summaries. // Keep the per-tier dimension in lockstep with the scalar counters; @@ -2108,6 +2382,32 @@ mod tests { assert_eq!(stats.tiers.get(UNKNOWN_TIER).map(|v| v.total_size), Some(9)); } + #[test] + fn unknown_tier_stats_deserialization_keeps_diagnostics_bounded() { + #[derive(Serialize)] + struct RawUnknownTierStats { + unknown_bytes: u64, + diagnostics: Vec, + diagnostics_dropped: u64, + } + + let mut diagnostics = vec!["tier-hash:0123456789abcdef".to_string(); UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP + 8]; + diagnostics.push("raw-tier-name-that-must-not-be-exposed".to_string()); + diagnostics.push("x".repeat(UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP + 1)); + let encoded = rmp_serde::to_vec_named(&RawUnknownTierStats { + unknown_bytes: 7, + diagnostics, + diagnostics_dropped: 3, + }) + .expect("unknown tier stats should encode"); + let decoded: UnknownTierStats = rmp_serde::from_slice(&encoded).expect("unknown tier stats should decode"); + + assert_eq!(decoded.unknown_bytes, 7); + assert_eq!(decoded.diagnostics.len(), 1); + assert!(decoded.diagnostics_dropped >= 3); + assert!(decoded.diagnostics.iter().all(|entry| is_redacted_tier_diagnostic(entry))); + } + #[test] fn checked_merge_rejects_overflowing_tier_totals() { let mut left = tier_entry( @@ -2131,6 +2431,57 @@ mod tests { assert_eq!(left.all_tier_stats.expect("left is untouched").tiers["WARM"].total_size, u64::MAX); } + #[test] + fn checked_merge_rejects_unknown_counter_overflow_on_either_side() { + let overflowed = DataUsageEntry { + unknown_tier_stats: Some(UnknownTierStats { + counter_overflowed: true, + ..Default::default() + }), + tier_accounting_proof: Some(TierAccountingProof { + logical_total: 1, + logical_known: 1, + ..Default::default() + }), + ..Default::default() + }; + let mut left = DataUsageEntry::default(); + assert!(!left.checked_merge(&overflowed)); + + let mut left = overflowed.clone(); + assert!(!left.checked_merge(&DataUsageEntry::default())); + } + + #[test] + fn checked_merge_rejects_one_sided_overflowed_proof_without_mutation() { + let mut left = DataUsageEntry { + size: 1, + tier_accounting_proof: Some(TierAccountingProof { + logical_total: 1, + logical_known: 1, + overflowed: true, + ..Default::default() + }), + ..Default::default() + }; + let legacy = DataUsageEntry { + size: 2, + ..Default::default() + }; + + assert!(!left.checked_merge(&legacy)); + assert_eq!(left.size, 1); + assert_eq!( + left.tier_accounting_proof, + Some(TierAccountingProof { + logical_total: 1, + logical_known: 1, + overflowed: true, + ..Default::default() + }) + ); + } + /// Entry shape released before per-tier accounting, using the derived /// (array) encoding those writers produced. #[derive(Serialize, Deserialize)] @@ -2170,6 +2521,72 @@ mod tests { assert_eq!(legacy.objects, 0); } + #[test] + fn tier_accounting_proof_is_optional_and_map_encoded() { + let entry = DataUsageEntry { + tier_accounting_proof: Some(TierAccountingProof { + logical_total: 11, + logical_known: 11, + physical_total: 7, + physical_known: 7, + ..Default::default() + }), + ..Default::default() + }; + let encoded = rmp_serde::to_vec(&entry).expect("proof-bearing entry should encode"); + let decoded: DataUsageEntry = rmp_serde::from_slice(&encoded).expect("proof-bearing entry should decode"); + assert_eq!(decoded.tier_accounting_proof, entry.tier_accounting_proof); + + let legacy = LegacyEntry { + children: DataUsageHashMap::default(), + size: 12, + objects: 3, + versions: 4, + delete_markers: 1, + obj_sizes: SizeHistogram::default(), + obj_versions: VersionsHistogram::default(), + replication_stats: None, + compacted: false, + failed_objects: 2, + }; + let legacy_bytes = rmp_serde::to_vec(&legacy).expect("legacy entry should encode"); + let decoded: DataUsageEntry = rmp_serde::from_slice(&legacy_bytes).expect("legacy entry should decode"); + assert!(decoded.tier_accounting_proof.is_none()); + } + + #[test] + fn empty_entry_merge_preserves_accounting_proof() { + let mut entry = DataUsageEntry { + tier_accounting_proof: Some(TierAccountingProof { + logical_total: 7, + logical_known: 7, + physical_total: 7, + physical_known: 7, + ..Default::default() + }), + ..Default::default() + }; + entry.merge(&DataUsageEntry::default()); + assert!(entry.tier_accounting_proof.is_some()); + + let mut empty = DataUsageEntry::default(); + empty.merge(&entry); + assert_eq!(empty.tier_accounting_proof, entry.tier_accounting_proof); + } + + #[test] + fn unknown_counter_overflow_is_not_empty_and_survives_roundtrip() { + let stats = UnknownTierStats { + counter_overflowed: true, + ..Default::default() + }; + assert!(!stats.is_empty()); + let encoded = rmp_serde::to_vec_named(&stats).expect("overflow marker should encode"); + let decoded: UnknownTierStats = rmp_serde::from_slice(&encoded).expect("overflow marker should decode"); + assert!(decoded.counter_overflowed); + assert!(!decoded.is_empty()); + } + #[test] fn legacy_array_encoded_entries_still_load() { let legacy = LegacyEntry { @@ -2938,6 +3355,32 @@ mod tests { ); } + #[test] + fn test_dui_filters_empty_unknown_tier_usage_from_the_flattened_tree() { + let root_hash = hash_path("root"); + let bucket_hash = hash_path("bucket-a"); + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: "root".to_string(), + ..Default::default() + }, + ..Default::default() + }; + cache.replace_hashed(&root_hash, &None, &DataUsageEntry::default()); + + let child = DataUsageEntry { + objects: 1, + unknown_tier_stats: Some(UnknownTierStats::default()), + ..Default::default() + }; + cache.replace_hashed(&bucket_hash, &Some(root_hash), &child); + + let info = cache.dui("root", &["bucket-a".to_string()]); + + assert_eq!(info.objects_total_count, 1); + assert!(info.unknown_tier_stats.is_none()); + } + #[test] fn test_data_usage_entry_merge_preserves_replication_targets() { let mut base = DataUsageEntry { diff --git a/crates/ecstore/src/api/mod.rs b/crates/ecstore/src/api/mod.rs index 456d38ac7..913b163dc 100644 --- a/crates/ecstore/src/api/mod.rs +++ b/crates/ecstore/src/api/mod.rs @@ -456,11 +456,13 @@ pub mod rpc { ScannerBucketListing, ScannerPeerActivity, TONIC_RPC_PREFIX, TonicInterceptor, build_put_file_auth_trailer, check_and_record_signed_rpc_nonce, gen_signature_headers, gen_tonic_replay_scope_headers, gen_tonic_signature_headers, gen_tonic_signature_interceptor, node_service_time_out_client, node_service_time_out_client_no_auth, - normalize_tonic_rpc_audience, set_tonic_canonical_body_digest, sign_ns_scanner_capability, sign_put_file_capability, - sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers, - tonic_rpc_auth_failure_reason, verify_put_file_auth_trailer, verify_put_file_capability, verify_rpc_signature, - verify_tonic_boot_epoch_response, verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest, - verify_tonic_rpc_response_proof, verify_tonic_rpc_signature, verify_tonic_rpc_signature_with_bootstrap, + normalize_tonic_rpc_audience, set_tonic_canonical_body_digest, sign_ns_scanner_capability, + sign_ns_scanner_capability_with_tier_registry_generation, sign_put_file_capability, sign_tonic_rpc_response_proof, + tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason, + verify_ns_scanner_capability, verify_ns_scanner_capability_with_tier_registry_generation, verify_put_file_auth_trailer, + verify_put_file_capability, verify_rpc_signature, verify_tonic_boot_epoch_response, verify_tonic_canonical_body_digest, + verify_tonic_mutation_body_digest, verify_tonic_rpc_response_proof, verify_tonic_rpc_signature, + verify_tonic_rpc_signature_with_bootstrap, }; } diff --git a/crates/ecstore/src/cluster/rpc/http_auth.rs b/crates/ecstore/src/cluster/rpc/http_auth.rs index bf05522ed..6be0bab22 100644 --- a/crates/ecstore/src/cluster/rpc/http_auth.rs +++ b/crates/ecstore/src/cluster/rpc/http_auth.rs @@ -94,6 +94,7 @@ const REPLAY_CACHE_AUTO_MEMORY_PERCENT: u64 = 13; const REPLAY_CACHE_AUTO_RPC_RPS_PER_CPU: usize = 4096; const REPLAY_CACHE_AUTO_MAX_CAPACITY: usize = 33_554_432; const NS_SCANNER_CAPABILITY_AUTH_DOMAIN: &[u8] = b"rustfs-ns-scanner-capability-v3"; +const NS_SCANNER_TIER_REGISTRY_GENERATION_AUTH_DOMAIN: &[u8] = b"rustfs-ns-scanner-tier-registry-generation-v1"; pub const TONIC_RPC_PREFIX: &str = "/node_service.NodeService"; static INTERNODE_RPC_SIGNATURE_STRICT: LazyLock = LazyLock::new(|| { get_env_bool( @@ -636,40 +637,79 @@ pub fn verify_put_file_capability(challenge: Uuid, server_epoch: Uuid, version: .map_err(|_| std::io::Error::new(std::io::ErrorKind::PermissionDenied, "Invalid put_file capability proof")) } -fn update_ns_scanner_capability_mac(mac: &mut HmacSha256, challenge: Uuid, server_epoch: Uuid) { +fn update_ns_scanner_capability_mac( + mac: &mut HmacSha256, + challenge: Uuid, + server_epoch: Uuid, + supports_tier_registry_generation: bool, +) { mac.update(NS_SCANNER_CAPABILITY_AUTH_DOMAIN); mac.update(&NS_SCANNER_PROTOCOL_VERSION.to_be_bytes()); mac.update(challenge.as_bytes()); mac.update(server_epoch.as_bytes()); + if supports_tier_registry_generation { + // The optional response capability is part of the authenticated + // scope. A proxy cannot turn an old/unsupported peer into a worker + // that receives generation-fenced scanner work. + mac.update(NS_SCANNER_TIER_REGISTRY_GENERATION_AUTH_DOMAIN); + } } -fn generate_ns_scanner_capability_proof(secret: &str, challenge: Uuid, server_epoch: Uuid) -> std::io::Result> { +fn generate_ns_scanner_capability_proof( + secret: &str, + challenge: Uuid, + server_epoch: Uuid, + supports_tier_registry_generation: bool, +) -> std::io::Result> { if challenge.is_nil() || server_epoch.is_nil() { return Err(std::io::Error::other("Invalid namespace scanner capability scope")); } let mut mac = ::new_from_slice(secret.as_bytes()).map_err(|_| std::io::Error::other("Invalid RPC HMAC key"))?; - update_ns_scanner_capability_mac(&mut mac, challenge, server_epoch); + update_ns_scanner_capability_mac(&mut mac, challenge, server_epoch, supports_tier_registry_generation); Ok(mac.finalize().into_bytes().to_vec()) } -fn verify_ns_scanner_capability_proof(secret: &str, challenge: Uuid, server_epoch: Uuid, proof: &[u8]) -> std::io::Result<()> { +fn verify_ns_scanner_capability_proof( + secret: &str, + challenge: Uuid, + server_epoch: Uuid, + proof: &[u8], + supports_tier_registry_generation: bool, +) -> std::io::Result<()> { if challenge.is_nil() || server_epoch.is_nil() { return Err(std::io::Error::other("Invalid namespace scanner capability scope")); } let mut mac = ::new_from_slice(secret.as_bytes()).map_err(|_| std::io::Error::other("Invalid RPC HMAC key"))?; - update_ns_scanner_capability_mac(&mut mac, challenge, server_epoch); + update_ns_scanner_capability_mac(&mut mac, challenge, server_epoch, supports_tier_registry_generation); mac.verify_slice(proof) .map_err(|_| std::io::Error::new(std::io::ErrorKind::PermissionDenied, "Invalid namespace scanner capability proof")) } pub fn sign_ns_scanner_capability(challenge: Uuid, server_epoch: Uuid) -> std::io::Result> { - generate_ns_scanner_capability_proof(&get_shared_secret()?, challenge, server_epoch) + sign_ns_scanner_capability_with_tier_registry_generation(challenge, server_epoch, false) } pub fn verify_ns_scanner_capability(challenge: Uuid, server_epoch: Uuid, proof: &[u8]) -> std::io::Result<()> { - verify_ns_scanner_capability_proof(&get_shared_secret()?, challenge, server_epoch, proof) + verify_ns_scanner_capability_with_tier_registry_generation(challenge, server_epoch, proof, false) +} + +pub fn sign_ns_scanner_capability_with_tier_registry_generation( + challenge: Uuid, + server_epoch: Uuid, + supports_tier_registry_generation: bool, +) -> std::io::Result> { + generate_ns_scanner_capability_proof(&get_shared_secret()?, challenge, server_epoch, supports_tier_registry_generation) +} + +pub fn verify_ns_scanner_capability_with_tier_registry_generation( + challenge: Uuid, + server_epoch: Uuid, + proof: &[u8], + supports_tier_registry_generation: bool, +) -> std::io::Result<()> { + verify_ns_scanner_capability_proof(&get_shared_secret()?, challenge, server_epoch, proof, supports_tier_registry_generation) } #[derive(Clone, Copy)] @@ -1709,13 +1749,28 @@ mod tests { let secret = "test-scanner-capability-secret"; let challenge = Uuid::new_v4(); let server_epoch = Uuid::new_v4(); - let proof = - generate_ns_scanner_capability_proof(secret, challenge, server_epoch).expect("capability proof should be generated"); + let proof = generate_ns_scanner_capability_proof(secret, challenge, server_epoch, false) + .expect("capability proof should be generated"); - assert!(verify_ns_scanner_capability_proof(secret, challenge, server_epoch, &proof).is_ok()); - assert!(verify_ns_scanner_capability_proof(secret, Uuid::new_v4(), server_epoch, &proof).is_err()); - assert!(verify_ns_scanner_capability_proof(secret, challenge, Uuid::new_v4(), &proof).is_err()); - assert!(verify_ns_scanner_capability_proof("different-secret", challenge, server_epoch, &proof).is_err()); + assert!(verify_ns_scanner_capability_proof(secret, challenge, server_epoch, &proof, false).is_ok()); + assert!(verify_ns_scanner_capability_proof(secret, Uuid::new_v4(), server_epoch, &proof, false).is_err()); + assert!(verify_ns_scanner_capability_proof(secret, challenge, Uuid::new_v4(), &proof, false).is_err()); + assert!(verify_ns_scanner_capability_proof("different-secret", challenge, server_epoch, &proof, false).is_err()); + } + + #[test] + fn namespace_scanner_capability_proof_binds_tier_registry_generation_support() { + let secret = "test-scanner-capability-secret"; + let challenge = Uuid::new_v4(); + let server_epoch = Uuid::new_v4(); + let proof = generate_ns_scanner_capability_proof(secret, challenge, server_epoch, true) + .expect("generation capability proof should be generated"); + + assert!(verify_ns_scanner_capability_proof(secret, challenge, server_epoch, &proof, true).is_ok()); + assert!(verify_ns_scanner_capability_proof(secret, challenge, server_epoch, &proof, false).is_err()); + let legacy = generate_ns_scanner_capability_proof(secret, challenge, server_epoch, false) + .expect("legacy capability proof should be generated"); + assert!(verify_ns_scanner_capability_proof(secret, challenge, server_epoch, &legacy, true).is_err()); } /// Security regression for GHSA-r5qv-rc46-hv8q (internode RPC fail-closed, diff --git a/crates/ecstore/src/cluster/rpc/internode_data_transport.rs b/crates/ecstore/src/cluster/rpc/internode_data_transport.rs index 6f46fead3..6b695954f 100644 --- a/crates/ecstore/src/cluster/rpc/internode_data_transport.rs +++ b/crates/ecstore/src/cluster/rpc/internode_data_transport.rs @@ -13,17 +13,18 @@ // limitations under the License. use crate::cluster::rpc::{ - build_auth_headers, build_put_file_auth_trailer, verify_ns_scanner_capability, verify_put_file_capability, + build_auth_headers, build_put_file_auth_trailer, verify_ns_scanner_capability_with_tier_registry_generation, + verify_put_file_capability, }; use crate::disk::error::{Error, Result}; use crate::disk::{FileReader, FileWriter}; use crate::storage_api_contracts::internode::{ NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY, NS_SCANNER_LEADER_EPOCH_QUERY, NS_SCANNER_PROTOCOL_VERSION, NS_SCANNER_PROTOCOL_VERSION_QUERY, NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY, - NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, NsScannerCapabilityResponse, PUT_FILE_AUTH_QUERY, - PUT_FILE_AUTH_V1, PUT_FILE_CAPABILITY_CHALLENGE_QUERY, PUT_FILE_CAPABILITY_QUERY, PUT_FILE_CAPABILITY_VERSION, - PUT_FILE_NONCE_QUERY, PUT_FILE_SERVER_EPOCH_QUERY, PutFileCapabilityResponse, WALK_DIR_BODY_SHA256_QUERY, - WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1, + NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY, + NsScannerCapabilityResponse, PUT_FILE_AUTH_QUERY, PUT_FILE_AUTH_V1, PUT_FILE_CAPABILITY_CHALLENGE_QUERY, + PUT_FILE_CAPABILITY_QUERY, PUT_FILE_CAPABILITY_VERSION, PUT_FILE_NONCE_QUERY, PUT_FILE_SERVER_EPOCH_QUERY, + PutFileCapabilityResponse, WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1, }; use async_trait::async_trait; use http::{HeaderMap, HeaderValue, Method, header::CONTENT_TYPE}; @@ -137,6 +138,12 @@ fn put_file_capability_status_is_legacy(status: u16) -> bool { status == 404 } +fn ns_scanner_capability_error_allows_legacy(error: &Error) -> bool { + [400, 404, 405, 426] + .into_iter() + .any(|status| error.is_internode_http_status(status)) +} + #[derive(Debug, Clone, Copy, Eq, PartialEq)] #[allow( dead_code, @@ -220,6 +227,7 @@ pub struct NsScannerStreamRequest { #[derive(Debug, Clone)] pub struct NsScannerCapabilityRequest { pub endpoint: String, + pub supports_tier_registry_generation: bool, } /// Data-plane stream opener used by `RemoteDisk`. @@ -252,6 +260,15 @@ pub trait InternodeDataTransport: Send + Sync + std::fmt::Debug { async fn probe_ns_scanner(&self, _request: NsScannerCapabilityRequest) -> Result { Err(Error::MethodNotAllowed) } + async fn probe_ns_scanner_capability(&self, request: NsScannerCapabilityRequest) -> Result { + let server_epoch = self.probe_ns_scanner(request).await?; + Ok(NsScannerCapabilityResponse { + version: NS_SCANNER_PROTOCOL_VERSION, + server_epoch, + proof: Vec::new(), + supports_tier_registry_generation: None, + }) + } // Interface facet nobody calls yet: every transport implements both, but no // caller negotiates on them. Kept for the internode transport split // (backlog#1350); deleting them would delete the seam and six impls. @@ -335,27 +352,44 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport { } async fn probe_ns_scanner(&self, request: NsScannerCapabilityRequest) -> Result { - let challenge = Uuid::new_v4(); - let url = build_ns_scanner_capability_url(&request, challenge); - let mut headers = msgpack_headers(); - build_auth_headers(&url, &Method::GET, &mut headers)?; - let reader = HttpReader::new(url, Method::GET, headers, None).await?; - let mut body = Vec::new(); - reader - .take(u64::try_from(NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE + 1).unwrap_or(u64::MAX)) - .read_to_end(&mut body) - .await?; - if body.is_empty() || body.len() > NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE { - return Err(Error::other("invalid remote namespace scanner capability response size")); + Ok(self.probe_ns_scanner_capability(request).await?.server_epoch) + } + + async fn probe_ns_scanner_capability(&self, request: NsScannerCapabilityRequest) -> Result { + if request.supports_tier_registry_generation { + return match self.probe_ns_scanner_capability_once(&request).await { + Ok(response) => Ok(response), + Err(marked_error) if ns_scanner_capability_error_allows_legacy(&marked_error) => { + // A v3 peer may reject the additive query marker, ignore + // it, or return its legacy proof. Retry once without the + // marker and only downgrade after that legacy response is + // authenticated; an unverified epoch is never trusted. + let legacy_request = NsScannerCapabilityRequest { + endpoint: request.endpoint.clone(), + supports_tier_registry_generation: false, + }; + match self.probe_ns_scanner_capability_once(&legacy_request).await { + Ok(mut response) => { + response.supports_tier_registry_generation = None; + Ok(response) + } + Err(legacy_error) if ns_scanner_capability_error_allows_legacy(&legacy_error) => { + // Some old deployments expose only the legacy + // protocol response (or advertise 426). Treat + // the pair as an explicit unsupported result so + // the scanner can use its coordinator fallback. + Err(Error::MethodNotAllowed) + } + Err(_) => Err(marked_error), + } + } + // A server failure, network failure, or authentication error + // is not evidence of an old parser. Do not issue an + // unauthenticated legacy probe or silently downgrade. + Err(marked_error) => Err(marked_error), + }; } - let response: NsScannerCapabilityResponse = - rmp_serde::from_slice(&body).map_err(|_| Error::other("invalid remote namespace scanner capability response"))?; - if response.version != NS_SCANNER_PROTOCOL_VERSION || response.server_epoch.is_nil() { - return Err(Error::other("incompatible remote namespace scanner capability response")); - } - verify_ns_scanner_capability(challenge, response.server_epoch, &response.proof) - .map_err(|err| Error::other(format!("remote namespace scanner capability authentication failed: {err}")))?; - Ok(response.server_epoch) + self.probe_ns_scanner_capability_once(&request).await } fn name(&self) -> &'static str { @@ -368,6 +402,53 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport { } impl TcpHttpInternodeDataTransport { + async fn probe_ns_scanner_capability_once( + &self, + request: &NsScannerCapabilityRequest, + ) -> Result { + let challenge = Uuid::new_v4(); + let url = build_ns_scanner_capability_url(request, challenge); + let mut headers = msgpack_headers(); + build_auth_headers(&url, &Method::GET, &mut headers)?; + let reader = HttpReader::new(url, Method::GET, headers, None).await?; + let mut body = Vec::new(); + reader + .take(u64::try_from(NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE + 1).unwrap_or(u64::MAX)) + .read_to_end(&mut body) + .await?; + if body.is_empty() || body.len() > NS_SCANNER_MAX_CAPABILITY_RESPONSE_SIZE { + return Err(Error::other("invalid remote namespace scanner capability response size")); + } + let mut response: NsScannerCapabilityResponse = + rmp_serde::from_slice(&body).map_err(|_| Error::other("invalid remote namespace scanner capability response"))?; + if response.version != NS_SCANNER_PROTOCOL_VERSION || response.server_epoch.is_nil() { + return Err(Error::other("incompatible remote namespace scanner capability response")); + } + if let Err(err) = verify_ns_scanner_capability_with_tier_registry_generation( + challenge, + response.server_epoch, + &response.proof, + request.supports_tier_registry_generation, + ) { + // A permissive older peer can ignore the additive marker and + // return a valid legacy-scope proof with HTTP 200. Accept that + // response only after independently authenticating the legacy + // scope; all other verification failures remain fail-closed. + if request.supports_tier_registry_generation && ns_scanner_capability_legacy_proof_is_valid(challenge, &response) { + response.supports_tier_registry_generation = None; + return Ok(response); + } + return Err(Error::other(format!("remote namespace scanner capability authentication failed: {err}"))); + } + // The proof authenticates the requested capability scope, not the + // optional response field. Derive the client-facing bit from that + // verified scope so an intermediary cannot strip or rewrite the field + // and force a silent downgrade after a successful generation-bound + // handshake. + normalize_ns_scanner_capability_response(&mut response, request.supports_tier_registry_generation); + Ok(response) + } + async fn put_file_auth_capability(&self, endpoint: &str) -> Result> { resolve_put_file_auth_capability(endpoint, || async { tokio::time::timeout(PUT_FILE_CAPABILITY_PROBE_TIMEOUT, self.probe_put_file_auth(endpoint)) @@ -649,6 +730,14 @@ fn build_walk_dir_url(request: &WalkDirStreamRequest) -> String { ) } +fn normalize_ns_scanner_capability_response(response: &mut NsScannerCapabilityResponse, requested_generation_support: bool) { + response.supports_tier_registry_generation = requested_generation_support.then_some(true); +} + +fn ns_scanner_capability_legacy_proof_is_valid(challenge: Uuid, response: &NsScannerCapabilityResponse) -> bool { + verify_ns_scanner_capability_with_tier_registry_generation(challenge, response.server_epoch, &response.proof, false).is_ok() +} + fn build_ns_scanner_url(request: &NsScannerStreamRequest) -> String { let body_sha256 = hex_simd::encode_to_string(Sha256::digest(&request.body), hex_simd::AsciiCase::Lower); format!( @@ -675,13 +764,18 @@ fn build_ns_scanner_url(request: &NsScannerStreamRequest) -> String { fn build_ns_scanner_capability_url(request: &NsScannerCapabilityRequest, challenge: Uuid) -> String { format!( - "{}{}?{}={}&{}={}", + "{}{}?{}={}&{}={}{}", request.endpoint, NS_SCANNER_PATH, NS_SCANNER_PROTOCOL_VERSION_QUERY, NS_SCANNER_PROTOCOL_VERSION, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, - challenge + challenge, + if request.supports_tier_registry_generation { + format!("&{}=true", NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY) + } else { + String::new() + } ) } @@ -794,6 +888,7 @@ mod tests { let probe_err = transport .probe_ns_scanner(NsScannerCapabilityRequest { endpoint: "http://node1:9000".to_string(), + supports_tier_registry_generation: false, }) .await .expect_err("legacy transport should report namespace scanner as unsupported"); @@ -1387,6 +1482,7 @@ mod tests { let url = build_ns_scanner_capability_url( &NsScannerCapabilityRequest { endpoint: "http://node1:9000".to_string(), + supports_tier_registry_generation: false, }, challenge, ); @@ -1399,6 +1495,85 @@ mod tests { ); } + #[test] + fn ns_scanner_capability_url_marks_generation_support_only_when_requested() { + let challenge = Uuid::new_v4(); + let url = build_ns_scanner_capability_url( + &NsScannerCapabilityRequest { + endpoint: "http://node1:9000".to_string(), + supports_tier_registry_generation: true, + }, + challenge, + ); + + assert!(url.contains(&format!("&{}=true", NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY))); + } + + #[test] + fn ns_scanner_capability_legacy_fallback_requires_explicit_compatibility_status() { + for status in [400, 404, 405, 426] { + let error = Error::from(rustfs_rio::new_test_internode_http_io_error( + rustfs_rio::InternodeHttpErrorKind::HttpStatus(http::StatusCode::from_u16(status).expect("test status")), + )); + assert!( + ns_scanner_capability_error_allows_legacy(&error), + "status {status} should permit legacy retry" + ); + } + + let marked_server_error = Error::from(rustfs_rio::new_test_internode_http_io_error( + rustfs_rio::InternodeHttpErrorKind::HttpStatus(http::StatusCode::INTERNAL_SERVER_ERROR), + )); + let network_error = Error::from(rustfs_rio::new_test_internode_http_io_error( + rustfs_rio::InternodeHttpErrorKind::ConnectionRefused, + )); + let authentication_error = Error::other("remote namespace scanner capability authentication failed"); + assert!(!ns_scanner_capability_error_allows_legacy(&marked_server_error)); + assert!(!ns_scanner_capability_error_allows_legacy(&network_error)); + assert!(!ns_scanner_capability_error_allows_legacy(&authentication_error)); + } + + #[test] + fn authenticated_ns_scanner_capability_ignores_unprotected_response_bit() { + let mut response = NsScannerCapabilityResponse { + version: NS_SCANNER_PROTOCOL_VERSION, + server_epoch: Uuid::new_v4(), + proof: Vec::new(), + supports_tier_registry_generation: None, + }; + + normalize_ns_scanner_capability_response(&mut response, true); + assert_eq!(response.supports_tier_registry_generation, Some(true)); + + response.supports_tier_registry_generation = Some(false); + normalize_ns_scanner_capability_response(&mut response, false); + assert_eq!(response.supports_tier_registry_generation, None); + } + + #[test] + fn ns_scanner_capability_accepts_only_authenticated_legacy_scope_after_marker_mismatch() { + crate::runtime::sources::ensure_test_rpc_secret(); + let challenge = Uuid::new_v4(); + let response = NsScannerCapabilityResponse { + version: NS_SCANNER_PROTOCOL_VERSION, + server_epoch: Uuid::new_v4(), + proof: crate::cluster::rpc::sign_ns_scanner_capability(challenge, Uuid::new_v4()) + .expect("placeholder proof should be generated"), + supports_tier_registry_generation: None, + }; + // A proof bound to a different challenge cannot authorize the legacy + // fallback, even though the response has the expected shape. + assert!(!ns_scanner_capability_legacy_proof_is_valid(challenge, &response)); + + let server_epoch = response.server_epoch; + let valid_response = NsScannerCapabilityResponse { + proof: crate::cluster::rpc::sign_ns_scanner_capability(challenge, server_epoch) + .expect("legacy proof should be generated"), + ..response + }; + assert!(ns_scanner_capability_legacy_proof_is_valid(challenge, &valid_response)); + } + #[test] fn transport_config_defaults_to_tcp_http() { let transport = build_internode_data_transport(None).unwrap(); diff --git a/crates/ecstore/src/cluster/rpc/mod.rs b/crates/ecstore/src/cluster/rpc/mod.rs index 9af454ac4..c8e3024eb 100644 --- a/crates/ecstore/src/cluster/rpc/mod.rs +++ b/crates/ecstore/src/cluster/rpc/mod.rs @@ -35,8 +35,9 @@ pub use http_auth::{ TONIC_RPC_PREFIX, build_auth_headers, build_put_file_auth_trailer, check_and_record_signed_rpc_nonce, gen_signature_headers, gen_tonic_replay_scope_headers, gen_tonic_signature_headers, normalize_tonic_rpc_audience, set_tonic_canonical_body_digest, set_tonic_mutation_body_digest, set_tonic_rolling_canonical_body_digest, set_tonic_rolling_mutation_body_digest, - sign_ns_scanner_capability, sign_put_file_capability, sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, - tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason, verify_ns_scanner_capability, verify_put_file_auth_trailer, + sign_ns_scanner_capability, sign_ns_scanner_capability_with_tier_registry_generation, sign_put_file_capability, + sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason, + verify_ns_scanner_capability, verify_ns_scanner_capability_with_tier_registry_generation, verify_put_file_auth_trailer, verify_put_file_capability, verify_rpc_signature, verify_tonic_boot_epoch_response, verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest, verify_tonic_rpc_response_proof, verify_tonic_rpc_signature, verify_tonic_rpc_signature_with_bootstrap, diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index 32bf3ae57..43f5154f9 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -778,14 +778,16 @@ impl RemoteDisk { if self.health.is_faulty() { return Err(DiskError::FaultyDisk); } - let probe = self.data_transport.probe_ns_scanner(NsScannerCapabilityRequest { + let probe = self.data_transport.probe_ns_scanner_capability(NsScannerCapabilityRequest { endpoint: self.endpoint.grid_host(), + supports_tier_registry_generation: true, }); let result = timeout(NS_SCANNER_CAPABILITY_PROBE_TIMEOUT, probe) .await .map_err(|_| DiskError::other("remote namespace scanner capability probe timed out"))?; match result { - Ok(server_epoch) => Ok(Some(server_epoch)), + Ok(response) if response.supports_tier_registry_generation == Some(true) => Ok(Some(response.server_epoch)), + Ok(_) => Ok(None), // RUSTFS_COMPAT_TODO(ns-scanner-rpc-v3): old peers and legacy transports lack the authenticated startup-epoch handshake. Remove after every supported peer implements namespace scanner protocol v3. Err(DiskError::MethodNotAllowed) => Ok(None), Err(err) @@ -3974,10 +3976,21 @@ mod tests { NsScannerProbe(NsScannerCapabilityRequest), } - #[derive(Debug, Clone, Default)] + #[derive(Debug, Clone)] struct RecordingInternodeDataTransport { calls: Arc>>, ns_scanner_probe_status: Arc>>, + ns_scanner_generation_support: Arc>>, + } + + impl Default for RecordingInternodeDataTransport { + fn default() -> Self { + Self { + calls: Arc::default(), + ns_scanner_probe_status: Arc::default(), + ns_scanner_generation_support: Arc::new(StdMutex::new(Some(true))), + } + } } #[derive(Clone, Debug)] @@ -4197,6 +4210,15 @@ mod tests { Self { calls: Arc::default(), ns_scanner_probe_status: Arc::new(StdMutex::new(Some(status))), + ns_scanner_generation_support: Arc::new(StdMutex::new(Some(true))), + } + } + + fn with_ns_scanner_generation_support(support: Option) -> Self { + Self { + calls: Arc::default(), + ns_scanner_probe_status: Arc::default(), + ns_scanner_generation_support: Arc::new(StdMutex::new(support)), } } @@ -4862,6 +4884,23 @@ mod tests { Ok(Uuid::from_u128(1)) } + async fn probe_ns_scanner_capability( + &self, + request: NsScannerCapabilityRequest, + ) -> Result { + let server_epoch = self.probe_ns_scanner(request).await?; + let supports_tier_registry_generation = *self + .ns_scanner_generation_support + .lock() + .expect("namespace scanner generation support lock poisoned"); + Ok(crate::storage_api_contracts::internode::NsScannerCapabilityResponse { + version: crate::storage_api_contracts::internode::NS_SCANNER_PROTOCOL_VERSION, + server_epoch, + proof: Vec::new(), + supports_tier_registry_generation, + }) + } + fn name(&self) -> &'static str { "recording" } @@ -6674,6 +6713,22 @@ mod tests { } } + #[tokio::test] + async fn test_remote_disk_namespace_scanner_capability_falls_back_without_generation_support() { + for support in [None, Some(false)] { + let transport = RecordingInternodeDataTransport::with_ns_scanner_generation_support(support); + let remote_disk = new_remote_disk_with_transport(Arc::new(transport)).await; + + assert_eq!( + remote_disk + .ns_scanner_server_epoch() + .await + .expect("missing generation support should be classified as unsupported"), + None + ); + } + } + #[tokio::test] async fn test_remote_disk_namespace_scanner_capability_rejects_legacy_transport() { let remote_disk = new_remote_disk_with_transport(Arc::new(RetryingOpenReadInternodeDataTransport::default())).await; diff --git a/crates/ecstore/src/storage_api_contracts/mod.rs b/crates/ecstore/src/storage_api_contracts/mod.rs index 78a3f9c2d..30d42ef2c 100644 --- a/crates/ecstore/src/storage_api_contracts/mod.rs +++ b/crates/ecstore/src/storage_api_contracts/mod.rs @@ -27,12 +27,12 @@ pub(crate) mod internode { NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY, NS_SCANNER_LEADER_EPOCH_QUERY, NS_SCANNER_PROTOCOL_VERSION, NS_SCANNER_PROTOCOL_VERSION_QUERY, NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY, NS_SCANNER_SESSION_ID_QUERY, - NS_SCANNER_SESSION_SEQUENCE_QUERY, NsScannerCapabilityResponse, PUT_FILE_AUTH_QUERY, PUT_FILE_AUTH_TRAILER_DIGEST_LEN, - PUT_FILE_AUTH_TRAILER_LEN, PUT_FILE_AUTH_TRAILER_MAC_LEN, PUT_FILE_AUTH_TRAILER_MAGIC, PUT_FILE_AUTH_V1, - PUT_FILE_CAPABILITY_CHALLENGE_QUERY, PUT_FILE_CAPABILITY_QUERY, PUT_FILE_CAPABILITY_VERSION, PUT_FILE_NONCE_QUERY, - PUT_FILE_SERVER_EPOCH_QUERY, PutFileCapabilityResponse, SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION, - SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION, SCANNER_ACTIVITY_PROTOCOL_VERSION, WALK_DIR_BODY_SHA256_QUERY, - WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1, + NS_SCANNER_SESSION_SEQUENCE_QUERY, NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY, NsScannerCapabilityResponse, + PUT_FILE_AUTH_QUERY, PUT_FILE_AUTH_TRAILER_DIGEST_LEN, PUT_FILE_AUTH_TRAILER_LEN, PUT_FILE_AUTH_TRAILER_MAC_LEN, + PUT_FILE_AUTH_TRAILER_MAGIC, PUT_FILE_AUTH_V1, PUT_FILE_CAPABILITY_CHALLENGE_QUERY, PUT_FILE_CAPABILITY_QUERY, + PUT_FILE_CAPABILITY_VERSION, PUT_FILE_NONCE_QUERY, PUT_FILE_SERVER_EPOCH_QUERY, PutFileCapabilityResponse, + SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION, SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION, SCANNER_ACTIVITY_PROTOCOL_VERSION, + WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_STREAM_COMPLETION_QUERY, WALK_DIR_STREAM_COMPLETION_V1, }; } diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index f6fd2b3d1..a20e7683e 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -29,8 +29,8 @@ use rustfs_config::ENV_SCANNER_CACHE_SAVE_TIMEOUT_SECS; pub use rustfs_data_usage::{ AllTierStats, BucketTargetUsageInfo, BucketUsageInfo, DATA_USAGE_OBJECT_NAME, DATA_USAGE_OBSERVED_OBJECT_NAME, DataUsageEntry, DataUsageHash, DataUsageHashMap, DataUsageInfo, LEGACY_DATA_USAGE_OBJECT_NAME, PrefixUsageEntry, - PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierStats, UNKNOWN_TIER, UnknownTierStats, - hash_path, prefix_usage_in_cache, + PrefixUsageQuery, PrefixUsageSummary, ReplTargetSizeSummary, SizeSummary, TierAccountingProof, TierStats, UNKNOWN_TIER, + UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP, UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP, UnknownTierStats, hash_path, prefix_usage_in_cache, }; use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf}; use tokio::time::{Duration, Instant, sleep, timeout}; @@ -209,45 +209,76 @@ impl ScannerSizeSummaryExt for SizeSummary { let logical_size = size.max(0); let size = usize::try_from(logical_size).unwrap_or(usize::MAX); self.total_size = self.total_size.saturating_add(size); + let logical_bytes = u64::try_from(logical_size).unwrap_or(u64::MAX); + let physical_bytes = u64::try_from(oi.size.max(0)).unwrap_or(0); + let mut proof = TierAccountingProof { + logical_total: logical_bytes, + logical_known: 0, + physical_total: physical_bytes, + physical_known: 0, + overflowed: false, + }; if oi.transitioned_object.free_version { + proof.logical_known = logical_bytes; + proof.physical_known = physical_bytes; + self.tier_accounting_proof.saturating_add(proof); return; } - let mut tier = oi.storage_class.clone().unwrap_or_else(|| storageclass::STANDARD.to_string()); - if oi.transitioned_object.status == TRANSITION_COMPLETE { - tier = oi.transitioned_object.tier.clone(); - } + let tier = if oi.transitioned_object.status == TRANSITION_COMPLETE { + oi.transitioned_object.tier.as_str() + } else { + oi.storage_class.as_deref().unwrap_or(storageclass::STANDARD) + }; let builtin_tier = tier == storageclass::STANDARD || tier == storageclass::RRS; let tier_registry_is_empty = self.tier_stats.is_empty() || (self.tier_stats.len() == 1 && self.tier_stats.contains_key(UNKNOWN_TIER)); - let known_tier = tier != UNKNOWN_TIER && (builtin_tier || self.tier_stats.contains_key(&tier)); + let known_tier = tier != UNKNOWN_TIER && (builtin_tier || self.tier_stats.contains_key(tier)); // With no configured tier, retain the historical empty-map shape for // ordinary STANDARD/RRS objects. A non-built-in key is still an // observable unknown and must create only the fixed bucket. if tier_registry_is_empty && known_tier { + proof.logical_known = logical_bytes; + proof.physical_known = physical_bytes; + self.tier_accounting_proof.saturating_add(proof); return; } - let accounting_key = if known_tier { tier.clone() } else { UNKNOWN_TIER.to_string() }; - let tier_stats = self.tier_stats.entry(accounting_key).or_default(); - let physical_size = u64::try_from(oi.size.max(0)).unwrap_or(0); + // Configured tiers and the fixed bucket are normally seeded, so the + // hot path can mutate them without allocating a key for every object. + // The fallback inserts only when a legacy/no-config summary sees its + // first unknown key. + let tier_stats = if known_tier { + if let Some(stats) = self.tier_stats.get_mut(tier) { + stats + } else { + self.tier_stats.entry(tier.to_owned()).or_default() + } + } else if let Some(stats) = self.tier_stats.get_mut(UNKNOWN_TIER) { + stats + } else { + self.tier_stats.entry(UNKNOWN_TIER.to_string()).or_default() + }; *tier_stats = tier_stats.add(&TierStats { - total_size: physical_size, + total_size: physical_bytes, num_versions: 1, num_objects: u64::from(oi.is_latest), }); - if !known_tier { - self.unknown_tier_stats.record_dimensions( - &tier, - u64::try_from(logical_size).unwrap_or(u64::MAX), - physical_size, - 1, - u64::from(oi.is_latest), - ); + if known_tier { + proof.logical_known = logical_bytes; + proof.physical_known = physical_bytes; } + if !known_tier { + self.unknown_tier_stats + .record_dimensions(tier, logical_bytes, physical_bytes, 1, u64::from(oi.is_latest)); + if self.unknown_tier_stats.counter_overflowed { + proof.overflowed = true; + } + } + self.tier_accounting_proof.saturating_add(proof); } } @@ -295,6 +326,11 @@ pub struct DataUsageEntryInfo { pub name: String, pub parent: String, pub entry: DataUsageEntry, + /// Registry generation used to classify this root entry. Older remote + /// workers omit it; callers must reject that result when a frozen cycle + /// requires generation fencing. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub tier_registry_generation: Option, } #[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq, Hash)] @@ -381,7 +417,8 @@ 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(17))?; + let field_count = 16 + usize::from(self.tier_registry_generation.is_some()); + let mut state = serializer.serialize_map(Some(field_count))?; state.serialize_entry("name", &self.name)?; state.serialize_entry("next_cycle", &self.next_cycle)?; state.serialize_entry("leader_epoch", &self.leader_epoch)?; @@ -398,7 +435,9 @@ impl Serialize for DataUsageCacheInfo { 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.serialize_entry("tier_registry_generation", &self.tier_registry_generation)?; + if let Some(generation) = self.tier_registry_generation { + state.serialize_entry("tier_registry_generation", &generation)?; + } state.end() } } @@ -422,10 +461,51 @@ impl DataUsageCache { /// Reconcile tier keys loaded from an older cache against the registry /// frozen for this scan. New metadata is already routed through /// `UNKNOWN_TIER`; this pass handles retired keys that predate that rule. + /// Legacy `TierStats` carries physical bytes only, so this migration does + /// not manufacture a logical unknown-byte value from that physical total. pub(crate) fn fold_retired_tiers(&mut self, tier_names: &[String]) { + let known_tiers = tier_names.iter().map(String::as_str).collect::>(); for entry in self.cache.values_mut() { - if let Some(tiers) = entry.all_tier_stats.as_mut() { - tiers.fold_unknown_tiers(tier_names.iter().map(String::as_str)); + let Some(tiers) = entry.all_tier_stats.as_mut() else { continue }; + let existing_unknown = tiers.tiers.get(UNKNOWN_TIER).cloned().unwrap_or_default(); + let companion_present = entry.unknown_tier_stats.as_ref().is_some_and(|stats| !stats.is_empty()); + let migrate_existing_unknown = !companion_present; + let mut retired = TierStats::default(); + let mut retired_key_found = false; + if migrate_existing_unknown { + retired = retired.add(&existing_unknown); + } + for (tier, stats) in &tiers.tiers { + if tier != UNKNOWN_TIER + && tier != storageclass::STANDARD + && tier != storageclass::RRS + && !known_tiers.contains(tier.as_str()) + { + retired_key_found = true; + retired = retired.add(stats); + } + } + tiers.fold_unknown_tiers(tier_names.iter().map(String::as_str)); + if !retired.is_empty() && !companion_present { + entry.add_unknown_tier_stats(&UnknownTierStats { + // The legacy map stores physical bytes only. Logical + // bytes remain zero until a fresh object scan observes + // them under the current metadata format. + unknown_physical_bytes: retired.total_size, + unknown_objects: retired.num_objects, + unknown_versions: retired.num_versions, + ..Default::default() + }); + // The legacy tier map has no logical-byte dimension, so a + // proof that classified this retired key as known cannot be + // repaired safely. Mark it unvalidated and require a fresh + // scan rather than guessing a logical subtraction. + entry.tier_accounting_proof = None; + } else if retired_key_found { + // A nonempty companion has no provenance tying it to the + // retired map keys. Reject the mixed cache until a fresh scan + // reconciles the dimensions instead of double-counting them. + entry.tier_accounting_proof = None; } } } diff --git a/crates/scanner/src/data_usage_define/tests.rs b/crates/scanner/src/data_usage_define/tests.rs index 4316aff2a..8882c9051 100644 --- a/crates/scanner/src/data_usage_define/tests.rs +++ b/crates/scanner/src/data_usage_define/tests.rs @@ -710,6 +710,12 @@ fn unknown_tier_is_bounded_and_accounted() { assert_eq!(summary.unknown_tier_stats.unknown_bytes, 11); assert_eq!(summary.unknown_tier_stats.unknown_physical_bytes, 11); assert_eq!(summary.unknown_tier_stats.unknown_objects, 1); + assert_eq!(summary.tier_accounting_proof.logical_total, 11); + assert_eq!(summary.tier_accounting_proof.logical_known, 0); + assert_eq!(summary.tier_accounting_proof.physical_total, 11); + assert_eq!(summary.tier_accounting_proof.physical_known, 0); + assert!(summary.unknown_tier_stats.diagnostics.len() <= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP); + assert!(summary.unknown_tier_stats.diagnostics.iter().map(String::len).sum::() <= UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP); assert!( summary .unknown_tier_stats @@ -743,6 +749,10 @@ fn unknown_tier_is_accounted_when_no_remote_tier_is_configured() { ..Default::default() }; summary.actions_accounting(&standard, 4, 4); + assert_eq!(summary.tier_accounting_proof.logical_total, 13); + assert_eq!(summary.tier_accounting_proof.logical_known, 4); + assert_eq!(summary.tier_accounting_proof.physical_total, 3 + 4); + assert_eq!(summary.tier_accounting_proof.physical_known, 4); assert_eq!(summary.tier_stats.len(), 1, "built-ins preserve the no-tier map shape"); } @@ -762,6 +772,8 @@ fn million_unique_tier_keys_do_not_grow_stats_map() { assert_eq!(summary.tier_stats.len(), 2); assert_eq!(summary.tier_stats[UNKNOWN_TIER].total_size, 1_000_000); assert_eq!(summary.unknown_tier_stats.unknown_bytes, 1_000_000); + assert!(summary.unknown_tier_stats.diagnostics.len() <= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP); + assert!(summary.unknown_tier_stats.diagnostics.iter().map(String::len).sum::() <= UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP); } #[test] @@ -786,20 +798,85 @@ fn unknown_tier_never_triggers_transition() { fn removed_tier_survives_restart_as_unknown() { let mut summary = SizeSummary::new(); summary.tier_stats.insert("COLD".to_string(), TierStats::default()); + summary.tier_stats.insert( + "RETIRED".to_string(), + TierStats { + total_size: 5, + num_versions: 1, + num_objects: 1, + }, + ); let object = ObjectInfo { - storage_class: Some("retired-tier".to_string()), + storage_class: Some("COLD".to_string()), size: 5, ..Default::default() }; summary.actions_accounting(&object, 5, 5); let mut entry = DataUsageEntry::default(); entry.add_tier_sizes(&summary.tier_stats); - entry.add_unknown_tier_stats(&summary.unknown_tier_stats); + entry.add_unknown_tier_stats(&UnknownTierStats { + unknown_bytes: 2, + unknown_physical_bytes: 2, + unknown_objects: 1, + unknown_versions: 1, + ..Default::default() + }); let encoded = rmp_serde::to_vec(&entry).expect("entry should encode"); let restored: DataUsageEntry = rmp_serde::from_slice(&encoded).expect("entry should decode"); + assert_eq!(restored.unknown_tier_stats.as_ref().map(|stats| stats.unknown_bytes), Some(2)); + assert_eq!( + restored.all_tier_stats.as_ref().expect("tier stats persisted").tiers["RETIRED"].total_size, + 5 + ); - assert_eq!(restored.all_tier_stats.expect("tier stats persisted").tiers[UNKNOWN_TIER].total_size, 5); - assert_eq!(restored.unknown_tier_stats.expect("unknown stats persisted").unknown_bytes, 5); + let mut cache = DataUsageCache::default(); + cache.replace("bucket", "", restored); + cache.fold_retired_tiers(&["COLD".to_string()]); + let folded = cache.cache.get(&hash_path("bucket").key()).expect("folded cache entry"); + assert_eq!( + folded.all_tier_stats.as_ref().expect("tier stats persisted").tiers[UNKNOWN_TIER].total_size, + 5 + ); + assert_eq!(folded.unknown_tier_stats.as_ref().map(|stats| stats.unknown_bytes), Some(2)); +} + +#[test] +fn retired_tier_fold_is_idempotent_and_rejects_mixed_companion_provenance() { + let mut cache = DataUsageCache::default(); + let mut entry = DataUsageEntry { + all_tier_stats: Some(AllTierStats { + tiers: HashMap::from([( + "RETIRED".to_string(), + TierStats { + total_size: 5, + num_versions: 1, + num_objects: 1, + }, + )]), + }), + ..Default::default() + }; + entry.unknown_tier_stats = Some(UnknownTierStats { + unknown_physical_bytes: 5, + ..Default::default() + }); + entry.tier_accounting_proof = Some(TierAccountingProof { + physical_total: 5, + physical_known: 5, + ..Default::default() + }); + cache.replace("bucket", "", entry); + + cache.fold_retired_tiers(&["COLD".to_string()]); + let first = cache.cache.get(&hash_path("bucket").key()).expect("entry").clone(); + assert_eq!(first.all_tier_stats.as_ref().expect("tiers").tiers[UNKNOWN_TIER].total_size, 5); + assert_eq!(first.unknown_tier_stats.as_ref().expect("companion").unknown_physical_bytes, 5); + assert!(first.tier_accounting_proof.is_none(), "mixed provenance must not publish"); + + cache.fold_retired_tiers(&["COLD".to_string()]); + let second = cache.cache.get(&hash_path("bucket").key()).expect("entry"); + assert_eq!(second.all_tier_stats.as_ref().expect("tiers").tiers[UNKNOWN_TIER].total_size, 5); + assert_eq!(second.unknown_tier_stats.as_ref().expect("companion").unknown_physical_bytes, 5); } #[test] @@ -846,6 +923,16 @@ fn unknown_tier_counter_uses_checked_arithmetic() { unknown.record("overflow", 1, 1, 1); assert_eq!(unknown.unknown_bytes, u64::MAX); assert_eq!(unknown.unknown_objects, 1); + assert!(unknown.counter_overflowed); + assert!(unknown.checked_add(&UnknownTierStats::default()).is_none()); + assert!( + unknown + .checked_add(&UnknownTierStats { + unknown_bytes: 1, + ..Default::default() + }) + .is_none() + ); } #[test] @@ -1101,7 +1188,7 @@ const USAGE_CACHE_WIRE_FIXTURE: &[u8] = &[ 0xdc, 0x00, 0x20, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0x03, 0xb0, 0x63, 0x61, 0x63, 0x68, 0x65, 0x5f, 0x6b, 0x65, 0x79, 0x5f, 0x66, 0x6f, 0x72, 0x6d, 0x61, 0x74, 0x01, 0x81, 0xab, 0x77, 0x69, 0x72, 0x65, 0x2d, 0x62, 0x75, 0x63, - 0x6b, 0x65, 0x74, 0x8c, 0xa8, 0x63, 0x68, 0x69, 0x6c, 0x64, 0x72, 0x65, 0x6e, 0x90, 0xa4, 0x73, 0x69, 0x7a, 0x65, 0xcd, 0x10, + 0x6b, 0x65, 0x74, 0x8b, 0xa8, 0x63, 0x68, 0x69, 0x6c, 0x64, 0x72, 0x65, 0x6e, 0x90, 0xa4, 0x73, 0x69, 0x7a, 0x65, 0xcd, 0x10, 0x00, 0xa7, 0x6f, 0x62, 0x6a, 0x65, 0x63, 0x74, 0x73, 0x03, 0xa8, 0x76, 0x65, 0x72, 0x73, 0x69, 0x6f, 0x6e, 0x73, 0x05, 0xae, 0x64, 0x65, 0x6c, 0x65, 0x74, 0x65, 0x5f, 0x6d, 0x61, 0x72, 0x6b, 0x65, 0x72, 0x73, 0x01, 0xa9, 0x6f, 0x62, 0x6a, 0x5f, 0x73, 0x69, 0x7a, 0x65, 0x73, 0x9b, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0xac, 0x6f, 0x62, 0x6a, 0x5f, @@ -1109,8 +1196,7 @@ const USAGE_CACHE_WIRE_FIXTURE: &[u8] = &[ 0x69, 0x63, 0x61, 0x74, 0x69, 0x6f, 0x6e, 0x5f, 0x73, 0x74, 0x61, 0x74, 0x73, 0xc0, 0xa9, 0x63, 0x6f, 0x6d, 0x70, 0x61, 0x63, 0x74, 0x65, 0x64, 0xc3, 0xae, 0x66, 0x61, 0x69, 0x6c, 0x65, 0x64, 0x5f, 0x6f, 0x62, 0x6a, 0x65, 0x63, 0x74, 0x73, 0x02, 0xae, 0x61, 0x6c, 0x6c, 0x5f, 0x74, 0x69, 0x65, 0x72, 0x5f, 0x73, 0x74, 0x61, 0x74, 0x73, 0x91, 0x81, 0xa4, 0x57, 0x41, 0x52, 0x4d, - 0x93, 0xcd, 0x08, 0x00, 0x02, 0x01, 0xb2, 0x75, 0x6e, 0x6b, 0x6e, 0x6f, 0x77, 0x6e, 0x5f, 0x74, 0x69, 0x65, 0x72, 0x5f, 0x73, - 0x74, 0x61, 0x74, 0x73, 0xc0, + 0x93, 0xcd, 0x08, 0x00, 0x02, 0x01, ]; #[test] diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 6cdd605c5..6a6a73b75 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -24,8 +24,11 @@ use bytes::Bytes; use http::HeaderMap; use rustfs_config::server_config::{Config as ServerConfig, get_global_server_config as config_get_global_server_config}; +use sha2::{Digest as _, Sha256}; +use std::collections::HashMap; use std::path::PathBuf; use std::sync::Arc; +use std::sync::LazyLock; use std::sync::RwLock; use std::time::{Duration, Instant}; use storage_api::owner::{ @@ -410,6 +413,7 @@ pub(crate) fn resolve_scanner_server_config() -> Option { /// How long the scanner caches the runtime tier-name list before re-reading /// the tier configuration manager. const TIER_NAME_CACHE_TTL: Duration = Duration::from_secs(30); +const MAX_TIER_REGISTRY_NAME_BYTES: usize = 256; /// Process-wide TTL cache of runtime tier names. /// @@ -424,6 +428,10 @@ const TIER_NAME_CACHE_TTL: Duration = Duration::from_secs(30); /// by key in per-object accounting and disappears on the next refresh. static TIER_NAME_CACHE: RwLock> = RwLock::new(None); static TIER_REGISTRY_GENERATION: AtomicU64 = AtomicU64::new(0); +static TIER_CYCLE_SNAPSHOTS: LazyLock>> = + LazyLock::new(|| RwLock::new(HashMap::new())); +static TIER_ACTIVE_CYCLES: LazyLock>> = LazyLock::new(|| RwLock::new(HashMap::new())); +static TIER_NAME_REFRESH_LOCK: LazyLock> = LazyLock::new(|| tokio::sync::Mutex::new(())); /// Return one immutable registry snapshot for a scanner unit of work. pub(crate) async fn runtime_tier_registry() -> TierRegistrySnapshot { @@ -436,41 +444,184 @@ pub(crate) async fn runtime_tier_registry() -> TierRegistrySnapshot { } } - let tiers = ecstore_get_global_tier_config_mgr().read().await.list_tiers(); - let names: Arc<[String]> = tiers.iter().map(|tier| tier.name.clone()).collect::>().into(); + // Serialize refreshes so a slower read of the old config cannot overwrite + // a newer snapshot published by a concurrent caller. + let _refresh_guard = TIER_NAME_REFRESH_LOCK.lock().await; + { + let cached = TIER_NAME_CACHE.read().unwrap_or_else(|err| err.into_inner()).clone(); + if let Some((refreshed_at, snapshot)) = cached + && refreshed_at.elapsed() < TIER_NAME_CACHE_TTL + { + return snapshot; + } + } + let previous = TIER_NAME_CACHE .read() .unwrap_or_else(|err| err.into_inner()) .as_ref() .map(|(_, snapshot)| snapshot.clone()); - let snapshot = previous - .as_ref() - .map(|snapshot| { - let refreshed = snapshot.refreshed(Ok(Arc::clone(&names))); - TierRegistrySnapshot { - generation: TIER_REGISTRY_GENERATION.fetch_add(1, Ordering::Relaxed).saturating_add(1), - ..refreshed + let names = ecstore_get_global_tier_config_mgr() + .read() + .await + .list_tiers() + .into_iter() + .map(|tier| tier.name) + .collect::>(); + let snapshot = match validate_tier_registry_names(names) { + Ok(names) => { + let generation = next_tier_registry_generation(); + match previous { + Some(previous) => TierRegistrySnapshot { + generation, + ..previous.refreshed(Ok(names)) + }, + None => TierRegistrySnapshot { + generation, + ..TierRegistrySnapshot::initial(names) + }, } - }) - .unwrap_or_else(|| TierRegistrySnapshot { - generation: TIER_REGISTRY_GENERATION.fetch_add(1, Ordering::Relaxed).saturating_add(1), - ..TierRegistrySnapshot::initial(names) - }); + } + Err(()) => match previous { + Some(previous) => previous.refreshed(Err(())), + None => TierRegistrySnapshot { + generation: next_tier_registry_generation(), + names: Arc::new([]), + refresh_failed: true, + }, + }, + }; *TIER_NAME_CACHE.write().unwrap_or_else(|err| err.into_inner()) = Some((Instant::now(), snapshot.clone())); snapshot } +fn next_tier_registry_generation() -> u64 { + TIER_REGISTRY_GENERATION + .fetch_update(Ordering::AcqRel, Ordering::Relaxed, |current| Some(current.saturating_add(1))) + .unwrap_or(u64::MAX) +} + +fn validate_tier_registry_names(mut names: Vec) -> Result, ()> { + if names.iter().any(|name| { + name.is_empty() + || name.len() > MAX_TIER_REGISTRY_NAME_BYTES + || name.bytes().any(|byte| byte.is_ascii_control()) + || name == UNKNOWN_TIER + || name == storageclass::STANDARD + || name == storageclass::RRS + }) { + return Err(()); + } + names.sort_unstable(); + if names.windows(2).any(|pair| pair[0] == pair[1]) { + return Err(()); + } + Ok(names.into()) +} + /// Tier names currently registered in the tier configuration, cached for /// `TIER_NAME_CACHE_TTL`. pub(crate) async fn runtime_tier_names() -> Arc<[String]> { runtime_tier_registry().await.names } +/// Return the immutable tier registry for one scanner cycle/leader pair. +/// Different buckets and disks belonging to the same cycle share this entry, +/// so a TTL refresh cannot split one published cycle across generations. +pub(crate) async fn runtime_tier_registry_for_cycle(cycle: u64, leader_epoch: u64) -> TierRegistrySnapshot { + let key = (cycle, leader_epoch); + prune_inactive_tier_cycle_snapshots(cycle, leader_epoch); + { + let cached = TIER_CYCLE_SNAPSHOTS.read().unwrap_or_else(|err| err.into_inner()); + if let Some(snapshot) = cached.get(&key) { + return snapshot.clone(); + } + } + + let mut snapshot = runtime_tier_registry().await; + // The registry generation describes the configuration snapshot, not the + // scan that consumed it. Keep it stable across cycles so a healthy cache + // can be reused; cycle and leader fencing are carried separately by the + // cache metadata and scan plan. + snapshot.generation = tier_registry_generation(&snapshot.names); + let mut cached = TIER_CYCLE_SNAPSHOTS.write().unwrap_or_else(|err| err.into_inner()); + if let Some(existing) = cached.get(&key) { + return existing.clone(); + } + cached.insert(key, snapshot.clone()); + snapshot +} + +fn prune_inactive_tier_cycle_snapshots(cycle: u64, leader_epoch: u64) { + let active = TIER_ACTIVE_CYCLES.read().unwrap_or_else(|err| err.into_inner()); + let mut snapshots = TIER_CYCLE_SNAPSHOTS.write().unwrap_or_else(|err| err.into_inner()); + snapshots.retain(|(entry_cycle, entry_epoch), _| { + active.contains_key(&(*entry_cycle, *entry_epoch)) + || *entry_epoch > leader_epoch + || (*entry_epoch == leader_epoch && *entry_cycle >= cycle) + }); +} + +pub(crate) struct TierRegistryCycleGuard { + key: (u64, u64), +} + +impl Drop for TierRegistryCycleGuard { + fn drop(&mut self) { + let mut active = TIER_ACTIVE_CYCLES.write().unwrap_or_else(|err| err.into_inner()); + if let Some(count) = active.get_mut(&self.key) { + *count = count.saturating_sub(1); + if *count == 0 { + active.remove(&self.key); + } + } + } +} + +pub(crate) fn begin_tier_registry_cycle(cycle: u64, leader_epoch: u64) -> TierRegistryCycleGuard { + let mut active = TIER_ACTIVE_CYCLES.write().unwrap_or_else(|err| err.into_inner()); + let count = active.entry((cycle, leader_epoch)).or_default(); + *count = count.saturating_add(1); + TierRegistryCycleGuard { + key: (cycle, leader_epoch), + } +} + +fn tier_registry_generation(names: &[String]) -> u64 { + let mut hasher = Sha256::new(); + hasher.update(b"rustfs-tier-registry-v1"); + for name in names { + hasher.update(u64::try_from(name.len()).unwrap_or(u64::MAX).to_le_bytes()); + hasher.update(name.as_bytes()); + } + let digest = hasher.finalize(); + let mut prefix = [0_u8; 8]; + prefix.copy_from_slice(&digest[..8]); + u64::from_le_bytes(prefix) +} + +/// Drop cycle snapshots only after the scanner has finished publishing a +/// cycle. In-flight or retryable cycles must retain their original registry; +/// TTL/capacity eviction could make a later bucket in the same cycle refresh +/// to a different generation. +pub(crate) fn complete_tier_registry_cycle(cycle: u64, leader_epoch: u64) { + let active = TIER_ACTIVE_CYCLES.read().unwrap_or_else(|err| err.into_inner()); + let mut cached = TIER_CYCLE_SNAPSHOTS.write().unwrap_or_else(|err| err.into_inner()); + cached.retain(|(entry_cycle, entry_epoch), _| { + active.contains_key(&(*entry_cycle, *entry_epoch)) + || *entry_epoch > leader_epoch + || (*entry_epoch == leader_epoch && *entry_cycle > cycle) + }); +} + /// Test-only cache reset; the production cache has no invalidation hook /// because the TTL is its only refresh path. #[cfg(test)] fn reset_tier_name_cache_for_test() { *TIER_NAME_CACHE.write().unwrap_or_else(|err| err.into_inner()) = None; + TIER_ACTIVE_CYCLES.write().unwrap_or_else(|err| err.into_inner()).clear(); + TIER_CYCLE_SNAPSHOTS.write().unwrap_or_else(|err| err.into_inner()).clear(); + TIER_REGISTRY_GENERATION.store(0, Ordering::Relaxed); } pub(crate) async fn enqueue_runtime_free_version(oi: ScannerObjectInfo) { @@ -679,6 +830,68 @@ mod tests { assert!(Arc::ptr_eq(&first, &second)); } + #[tokio::test] + async fn tier_registry_cycle_snapshot_stays_fixed_while_active() { + reset_tier_name_cache_for_test(); + let cycle = 9_000_001; + let leader_epoch = 9_000_002; + let guard = begin_tier_registry_cycle(cycle, leader_epoch); + let first = runtime_tier_registry_for_cycle(cycle, leader_epoch).await; + + // Simulate a TTL refresh observing a different configuration while the + // original cycle is still scanning. The active cycle entry must win. + *TIER_NAME_CACHE.write().unwrap_or_else(|err| err.into_inner()) = Some(( + Instant::now() - TIER_NAME_CACHE_TTL - Duration::from_secs(1), + TierRegistrySnapshot { + generation: u64::MAX, + names: Arc::from(["COLD".to_string()]), + refresh_failed: false, + }, + )); + let second = runtime_tier_registry_for_cycle(cycle, leader_epoch).await; + assert_eq!(second.generation, first.generation); + assert_eq!(second.names, first.names); + + drop(guard); + complete_tier_registry_cycle(cycle, leader_epoch); + reset_tier_name_cache_for_test(); + } + + #[tokio::test] + async fn tier_registry_generation_survives_new_cycle_with_same_names() { + reset_tier_name_cache_for_test(); + let first_cycle = 9_000_011; + let second_cycle = first_cycle + 1; + let leader_epoch = 9_000_012; + + let first_guard = begin_tier_registry_cycle(first_cycle, leader_epoch); + let first = runtime_tier_registry_for_cycle(first_cycle, leader_epoch).await; + drop(first_guard); + complete_tier_registry_cycle(first_cycle, leader_epoch); + + let second_guard = begin_tier_registry_cycle(second_cycle, leader_epoch); + let second = runtime_tier_registry_for_cycle(second_cycle, leader_epoch).await; + assert_eq!(first.names, second.names); + assert_eq!(first.generation, second.generation); + + drop(second_guard); + complete_tier_registry_cycle(second_cycle, leader_epoch); + reset_tier_name_cache_for_test(); + } + + #[test] + fn invalid_tier_registry_names_fail_closed_for_refresh() { + assert!(validate_tier_registry_names(vec!["COLD\n".to_string()]).is_err()); + assert!(validate_tier_registry_names(vec![UNKNOWN_TIER.to_string()]).is_err()); + assert!(validate_tier_registry_names(vec!["COLD".to_string(), "COLD".to_string()]).is_err()); + assert_eq!( + validate_tier_registry_names(vec!["WARM".to_string(), "COLD".to_string()]) + .expect("valid registry names") + .as_ref(), + ["COLD".to_string(), "WARM".to_string()] + ); + } + #[test] fn foreground_read_guard_tracks_stream_lifetime() { reset_foreground_read_activity_for_test(); diff --git a/crates/scanner/src/remote_scanner/stream.rs b/crates/scanner/src/remote_scanner/stream.rs index 93ba33100..a6469d39e 100644 --- a/crates/scanner/src/remote_scanner/stream.rs +++ b/crates/scanner/src/remote_scanner/stream.rs @@ -16,8 +16,8 @@ use crate::RUSTFS_META_BUCKET; use crate::scanner_budget::{ScannerCycleBudget, ScannerCycleBudgetConfig}; use crate::scanner_io::{ - DataUsageCacheScanState, ScannerDiskScanOutcome, ScannerIODisk, acquire_scanner_cache_locks, cache_root_entry_info, - current_cache_root_or_prepare, scanner_set_disk_inventory, + DataUsageCacheReuseOptions, DataUsageCacheScanState, ScannerDiskScanOutcome, ScannerIODisk, acquire_scanner_cache_locks, + cache_root_entry_info, current_cache_root_or_prepare_with_generation, scanner_set_disk_inventory, }; use crate::storage_api::owner::NS_SCANNER_PROTOCOL_VERSION; use crate::{ @@ -212,6 +212,7 @@ pub(crate) struct RemoteScannerScanSpec<'a> { pub(crate) session_id: Uuid, pub(crate) session_sequence: u64, pub(crate) scan_plan_digest: DataUsageScanPlanDigest, + pub(crate) tier_registry_generation: u64, pub(crate) skip_healing: bool, pub(crate) scan_mode: HealScanMode, } @@ -222,6 +223,7 @@ struct RemoteScannerResponseExpectation<'a> { source: DataUsageCacheSource, next_cycle: u64, scan_plan_digest: DataUsageScanPlanDigest, + tier_registry_generation: u64, } #[derive(Debug)] @@ -668,6 +670,11 @@ async fn scan_and_persist_local_bucket( scan_mode, .. } = request; + // Keep the worker's cycle snapshot alive through cache reuse, scanning, + // and persistence. Without the guard, a later cycle can prune this key + // while this request is still running and allow a second registry to be + // selected for the same cycle. + let _tier_cycle_guard = crate::begin_tier_registry_cycle(next_cycle, leader_epoch); let store = resolve_scanner_object_store_handle() .ok_or_else(|| RemoteScannerServerError::worker("remote namespace scanner object layer is unavailable"))?; validate_remote_scanner_request_fence_with_store(next_cycle, leader_epoch, store.clone()) @@ -700,7 +707,25 @@ 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}")) })?; - let scan_state = current_cache_root_or_prepare(&mut cache, &bucket, source, next_cycle, leader_epoch, scan_plan_digest, true); + // Remote workers use the same cycle-frozen registry as `scan_data_folder`. + // Requiring its generation here prevents a cache snapshot classified by an + // older registry from being reused before the folder scan gets a chance to + // refresh it. + let tier_registry_generation = crate::runtime_tier_registry_for_cycle(next_cycle, leader_epoch) + .await + .generation; + let scan_state = current_cache_root_or_prepare_with_generation( + &mut cache, + &bucket, + source, + next_cycle, + leader_epoch, + scan_plan_digest, + DataUsageCacheReuseOptions { + require_source: true, + tier_registry_generation: Some(tier_registry_generation), + }, + ); match scan_state { DataUsageCacheScanState::Current(usage) => { if guard.is_lock_lost() { @@ -844,6 +869,7 @@ pub(crate) async fn scan_remote_bucket( session_id, session_sequence, scan_plan_digest, + tier_registry_generation, skip_healing, scan_mode, } = spec; @@ -932,6 +958,7 @@ pub(crate) async fn scan_remote_bucket( source: expected_source, next_cycle, scan_plan_digest, + tier_registry_generation, }, authenticator, rpc_deadline, @@ -987,6 +1014,7 @@ where source: expected_source, next_cycle: TEST_NEXT_CYCLE, scan_plan_digest: expected_scan_plan_digest, + tier_registry_generation: 0, }, authenticator, Instant::now() + NS_SCANNER_MAX_RPC_LIFETIME, @@ -1086,6 +1114,11 @@ where "remote namespace scanner returned usage for a different bucket plan", ))); } + if complete.usage.tier_registry_generation != Some(expected.tier_registry_generation) { + return Err(RemoteScannerStreamError::reconciled(StorageError::other( + "remote namespace scanner returned usage for a different tier registry generation", + ))); + } if !complete.usage.entry.children.is_empty() { return Err(RemoteScannerStreamError::reconciled(StorageError::other( "remote namespace scanner returned non-flattened bucket usage", diff --git a/crates/scanner/src/remote_scanner/stream/tests.rs b/crates/scanner/src/remote_scanner/stream/tests.rs index 11d28392e..ba0b71e58 100644 --- a/crates/scanner/src/remote_scanner/stream/tests.rs +++ b/crates/scanner/src/remote_scanner/stream/tests.rs @@ -195,6 +195,7 @@ fn test_usage(bucket: &str, objects: usize) -> DataUsageEntryInfo { name: bucket.to_string(), parent: crate::DATA_USAGE_ROOT.to_string(), entry, + tier_registry_generation: Some(0), } } @@ -751,6 +752,43 @@ async fn complete_terminal_frame_reconciles_progress_and_usage() { assert_eq!(budget.progress(), (3, 2)); } +#[tokio::test] +async fn terminal_usage_from_a_different_tier_generation_is_rejected() { + let request_id = Uuid::new_v4(); + let writer_auth = FrameAuthenticator::for_test(request_id); + let reader_auth = FrameAuthenticator::for_test(request_id); + let (mut writer, reader) = tokio::io::duplex(4096); + tokio::spawn(async move { + let mut usage = test_usage("bucket", 1); + usage.tier_registry_generation = Some(1); + let mut sequence = 0; + write_frame( + &mut writer, + &writer_auth, + &mut sequence, + &RemoteScannerFrame::terminal( + RemoteScannerProgress::default(), + RemoteScannerFrameResult::Complete(Box::new(RemoteScannerComplete { + source: TEST_SOURCE, + scan_plan_digest: TEST_PLAN_DIGEST, + usage, + pending_maintenance_work: false, + })), + ), + ) + .await + .expect("terminal frame should write"); + }); + + let parent = CancellationToken::new(); + let budget = ScannerCycleBudget::new(&parent, ScannerCycleBudgetConfig::default()); + let error = consume_remote_scanner_stream(reader, parent, budget, "bucket", TEST_SOURCE, TEST_PLAN_DIGEST, reader_auth) + .await + .expect_err("generation mismatch must fail closed"); + + assert!(error.to_string().contains("tier registry generation")); +} + #[tokio::test] async fn complete_terminal_frame_after_budget_expiry_is_partial() { let request_id = Uuid::new_v4(); diff --git a/crates/scanner/src/scanner_folder.rs b/crates/scanner/src/scanner_folder.rs index 28ebaf2b8..19a5b0a41 100644 --- a/crates/scanner/src/scanner_folder.rs +++ b/crates/scanner/src/scanner_folder.rs @@ -57,7 +57,7 @@ use crate::{ ReplicationHealObject, ReplicationQueueAdmission, ReplicationStatusType, STORAGE_FORMAT_FILE, ScannerDiskExt as _, ScannerLifecycleConfigExt as _, ScannerVersioningConfigExt as _, StorageError, TierRegistrySnapshot, apply_expiry_rule, apply_transition_rule, enqueue_runtime_newer_noncurrent, is_reserved_or_invalid_bucket, list_path_raw, path2_bucket_object, - path2_bucket_object_with_base_path, queue_replication_heal, runtime_tier_registry, scanner_is_erasure, + path2_bucket_object_with_base_path, queue_replication_heal, runtime_tier_registry_for_cycle, scanner_is_erasure, scanner_replication_config_for_lifecycle_eval, }; use crate::{ScannerObjectInfo as ObjectInfo, ScannerObjectToDelete as ObjectToDelete}; @@ -627,6 +627,19 @@ fn apply_scanner_size_summary(into: &mut DataUsageEntry, summary: &SizeSummary) into.add_tier_sizes(&summary.tier_stats); into.add_unknown_tier_stats(&summary.unknown_tier_stats); + into.tier_accounting_proof = match (into.tier_accounting_proof, Some(summary.tier_accounting_proof)) { + (Some(mut current), Some(next)) => { + current.saturating_add(next); + Some(current) + } + (Some(_), None) => None, + (None, next) => next, + }; + if into.unknown_tier_stats.as_ref().is_some_and(|stats| stats.counter_overflowed) + && let Some(proof) = into.tier_accounting_proof.as_mut() + { + proof.overflowed = true; + } } fn data_usage_root_has_progress(root: &DataUsageEntry) -> bool { @@ -637,6 +650,8 @@ fn data_usage_root_has_progress(root: &DataUsageEntry) -> bool { || root.delete_markers > 0 || root.failed_objects > 0 || root.replication_stats.is_some() + || root.all_tier_stats.as_ref().is_some_and(|stats| !stats.is_empty()) + || root.unknown_tier_stats.as_ref().is_some_and(|stats| !stats.is_empty()) } fn partial_cache_is_useful(root: &DataUsageEntry, pending_heals_changed: bool) -> bool { @@ -2169,7 +2184,7 @@ pub async fn scan_data_folder( let failed_object_ttl = rustfs_utils::get_env_u32(ENV_FAILED_OBJECT_TTL_SECS, DEFAULT_FAILED_OBJECT_TTL_SECS) as u64; let failed_objects_max = rustfs_utils::get_env_u32(ENV_FAILED_OBJECTS_MAX, DEFAULT_FAILED_OBJECTS_MAX) as usize; - let tier_registry = runtime_tier_registry().await; + let tier_registry = runtime_tier_registry_for_cycle(cache.info.next_cycle, cache.info.leader_epoch).await; let mut cache = cache; cache.fold_retired_tiers(&tier_registry.names); cache.info.tier_registry_generation = Some(tier_registry.generation); diff --git a/crates/scanner/src/scanner_folder/item_actions.rs b/crates/scanner/src/scanner_folder/item_actions.rs index 91fd54c27..3c29e26ed 100644 --- a/crates/scanner/src/scanner_folder/item_actions.rs +++ b/crates/scanner/src/scanner_folder/item_actions.rs @@ -297,12 +297,16 @@ impl ScannerItem { } } - fn tier_is_known(oi: &ObjectInfo, tier_names: &[String]) -> bool { - let tier = Self::effective_tier(oi); - if tier == crate::data_usage_define::UNKNOWN_TIER { - return false; - } - tier == crate::storageclass::STANDARD || tier == crate::storageclass::RRS || tier_names.iter().any(|name| name == tier) + fn tier_name_is_known(tier: &str, tier_names: &[String]) -> bool { + !tier.is_empty() + && tier != crate::data_usage_define::UNKNOWN_TIER + && (tier == crate::storageclass::STANDARD + || tier == crate::storageclass::RRS + || tier_names.iter().any(|name| name == tier)) + } + + pub(crate) fn tier_is_known(oi: &ObjectInfo, tier_names: &[String]) -> bool { + Self::tier_name_is_known(Self::effective_tier(oi), tier_names) } fn action_requires_known_tier(action: IlmAction) -> bool { @@ -319,6 +323,22 @@ impl ScannerItem { ) } + fn action_blocked_by_unknown_tier( + action: IlmAction, + oi: &ObjectInfo, + all_versions_known: bool, + tier_names: &[String], + target: &str, + ) -> bool { + if !Self::action_requires_known_tier(action) { + return false; + } + !Self::tier_is_known(oi, tier_names) + || (action.delete_all() && !all_versions_known) + || (matches!(action, IlmAction::TransitionAction | IlmAction::TransitionVersionAction) + && !Self::tier_name_is_known(target, tier_names)) + } + pub async fn apply_actions( &mut self, object_infos: Vec, @@ -433,6 +453,9 @@ impl ScannerItem { let mut noncurrent_accounting: Vec> = Vec::new(); let mut cumulative_size = 0; let mut remaining_versions = object_infos.len(); + let all_versions_known = object_infos + .iter() + .all(|candidate| Self::tier_is_known(candidate, tier_names)); 'eventLoop: { for (i, event) in events.iter().enumerate() { let oi = &object_infos[i]; @@ -461,7 +484,7 @@ impl ScannerItem { // evaluation is still useful for accounting, but all // side-effecting tier actions fail closed until the registry // recognizes the source again. - if !Self::tier_is_known(oi, tier_names) && Self::action_requires_known_tier(event.action) { + if Self::action_blocked_by_unknown_tier(event.action, oi, all_versions_known, tier_names, &event.storage_class) { size = self.heal_actions(oi, actual_size, size_summary).await; size_summary.actions_accounting(oi, size, actual_size); cumulative_size += size; @@ -979,8 +1002,42 @@ mod tests { storage_class: Some("retired-tier".to_string()), ..Default::default() }; - assert!(!ScannerItem::tier_is_known(&object, &["WARM".to_string()])); + let tier_names = ["WARM".to_string()]; + assert!(!ScannerItem::tier_is_known(&object, &tier_names)); assert!(ScannerItem::action_requires_known_tier(IlmAction::TransitionAction)); assert!(ScannerItem::action_requires_known_tier(IlmAction::DeleteVersionAction)); + assert!(ScannerItem::action_blocked_by_unknown_tier( + IlmAction::TransitionAction, + &object, + false, + &tier_names, + "WARM" + )); + assert!(!ScannerItem::action_blocked_by_unknown_tier( + IlmAction::NoneAction, + &object, + false, + &tier_names, + "WARM" + )); + + let known = ObjectInfo { + storage_class: Some(crate::storageclass::STANDARD.to_string()), + ..Default::default() + }; + assert!(ScannerItem::action_blocked_by_unknown_tier( + IlmAction::DeleteAllVersionsAction, + &known, + false, + &tier_names, + "WARM" + )); + assert!(!ScannerItem::action_blocked_by_unknown_tier( + IlmAction::TransitionAction, + &known, + true, + &tier_names, + crate::storageclass::STANDARD + )); } } diff --git a/crates/scanner/src/scanner_io.rs b/crates/scanner/src/scanner_io.rs index a8219e606..6ea694616 100644 --- a/crates/scanner/src/scanner_io.rs +++ b/crates/scanner/src/scanner_io.rs @@ -57,8 +57,9 @@ use crate::storage_api::scanner_io::{BucketInfo, BucketOptions}; use crate::{ BucketTargetSys, BucketVersioningSys, Disk, DiskError, ECStore, EcstoreError as Error, EcstoreResult as Result, RUSTFS_META_BUCKET, ReplicationConfig, STORAGE_FORMAT_FILE, ScannerDiskExt as _, ScannerLifecycleConfigExt as _, - ScannerReplicationConfigExt as _, ScannerVersioningConfigExt as _, SetDisks, StorageError, enqueue_runtime_free_version, - get_lifecycle_config, get_object_lock_config, get_replication_config, runtime_tier_names, storageclass, + ScannerReplicationConfigExt as _, ScannerVersioningConfigExt as _, SetDisks, StorageError, begin_tier_registry_cycle, + complete_tier_registry_cycle, enqueue_runtime_free_version, get_lifecycle_config, get_object_lock_config, + get_replication_config, runtime_tier_names, runtime_tier_registry_for_cycle, storageclass, }; pub(crate) const SCANNER_SKIP_FILE_ERROR: &str = "skip file"; @@ -143,6 +144,7 @@ pub struct ScannerBucketScanPlan { all_buckets: Arc>, digest: DataUsageScanPlanDigest, leader_epoch: u64, + tier_registry_generation: u64, dirty_usage_buckets: Arc, bucket_failures: ScannerBucketFailureState, pending_maintenance_work: Arc, @@ -341,12 +343,20 @@ pub(crate) fn cache_root_entry_info(cache: &DataUsageCache) -> std::result::Resu name: cache.info.name.clone(), parent: DATA_USAGE_ROOT.to_string(), entry, + tier_registry_generation: cache.info.tier_registry_generation, }) } -fn apply_bucket_result_to_cache(cache: &mut DataUsageCache, result: DataUsageEntryInfo, update_time: SystemTime) { +fn apply_bucket_result_to_cache(cache: &mut DataUsageCache, result: DataUsageEntryInfo, update_time: SystemTime) -> bool { + if cache.info.tier_registry_generation != result.tier_registry_generation { + // A result from another registry generation must never be folded into + // this cycle. Leaving it unapplied makes the cycle incomplete and + // forces the caller to re-account it under one frozen registry. + return false; + } cache.replace(&result.name, &result.parent, result.entry); cache.info.last_update = Some(update_time); + true } fn should_publish_completed_snapshot(completed_count: usize, total_count: usize, budget_elapsed: bool, cancelled: bool) -> bool { @@ -652,7 +662,10 @@ use cache::*; use dirty_usage::*; use guards::*; -pub(crate) use cache::{DataUsageCacheScanState, acquire_scanner_cache_locks, current_cache_root_or_prepare}; +pub(crate) use cache::{ + DataUsageCacheReuseOptions, DataUsageCacheScanState, acquire_scanner_cache_locks, + current_cache_root_or_prepare_with_generation, +}; pub use dirty_usage::{ ScannerDirtyUsageAckError, ScannerDirtyUsageState, acknowledge_dirty_usage_generation, clear_dirty_usage_bucket, record_dirty_usage_bucket, record_scanner_maintenance_change, scanner_activity_epoch, scanner_dirty_usage_state, diff --git a/crates/scanner/src/scanner_io/cache.rs b/crates/scanner/src/scanner_io/cache.rs index 2b64cad38..3ad8bdbb8 100644 --- a/crates/scanner/src/scanner_io/cache.rs +++ b/crates/scanner/src/scanner_io/cache.rs @@ -100,13 +100,14 @@ where let _ = tokio::time::timeout(SCANNER_CACHE_LOCK_LOSS_SHUTDOWN_TIMEOUT, scan).await; } -pub(crate) fn current_cache_root_entry( +pub(crate) fn current_cache_root_entry_with_generation( cache: &DataUsageCache, name: &str, source: DataUsageCacheSource, next_cycle: u64, leader_epoch: u64, scan_plan_digest: DataUsageScanPlanDigest, + tier_registry_generation: Option, ) -> std::result::Result, ScannerError> { let metadata_is_current = cache.info.name == name && cache.info.source == Some(source) @@ -115,7 +116,8 @@ pub(crate) fn current_cache_root_entry( && 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; + && cache.info.cache_key_format == DATA_USAGE_CACHE_KEY_FORMAT + && tier_registry_generation.is_none_or(|generation| cache.info.tier_registry_generation == Some(generation)); if !metadata_is_current { return Ok(None); } @@ -131,6 +133,13 @@ pub(crate) enum DataUsageCacheScanState { }, } +#[derive(Clone, Copy, Debug, Default)] +pub(crate) struct DataUsageCacheReuseOptions { + pub(crate) require_source: bool, + pub(crate) tier_registry_generation: Option, +} + +#[cfg(test)] pub(crate) fn current_cache_root_or_prepare( cache: &mut DataUsageCache, name: &str, @@ -140,11 +149,51 @@ pub(crate) fn current_cache_root_or_prepare( scan_plan_digest: DataUsageScanPlanDigest, require_source: bool, ) -> DataUsageCacheScanState { - match current_cache_root_entry(cache, name, source, next_cycle, leader_epoch, scan_plan_digest) { + current_cache_root_or_prepare_with_generation( + cache, + name, + source, + next_cycle, + leader_epoch, + scan_plan_digest, + DataUsageCacheReuseOptions { + require_source, + tier_registry_generation: None, + }, + ) +} + +pub(crate) fn current_cache_root_or_prepare_with_generation( + cache: &mut DataUsageCache, + name: &str, + source: DataUsageCacheSource, + next_cycle: u64, + leader_epoch: u64, + scan_plan_digest: DataUsageScanPlanDigest, + options: DataUsageCacheReuseOptions, +) -> DataUsageCacheScanState { + if options.tier_registry_generation.is_some_and(|generation| { + cache.info.next_cycle <= next_cycle + && cache.info.leader_epoch <= leader_epoch + && cache.info.tier_registry_generation != Some(generation) + }) { + // Make prepare_for_scan take its reset path so an entry classified by + // an older registry cannot be reused under the new cycle generation. + cache.info.scan_plan_digest = None; + } + match current_cache_root_entry_with_generation( + cache, + name, + source, + next_cycle, + leader_epoch, + scan_plan_digest, + options.tier_registry_generation, + ) { 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), + outcome: cache.prepare_for_scan(name, next_cycle, leader_epoch, source, scan_plan_digest, options.require_source), }, } } @@ -159,7 +208,7 @@ pub(super) fn cache_snapshot_is_current( scan_plan_digest: DataUsageScanPlanDigest, ) -> bool { matches!( - current_cache_root_entry(cache, name, source, next_cycle, leader_epoch, scan_plan_digest), + current_cache_root_entry_with_generation(cache, name, source, next_cycle, leader_epoch, scan_plan_digest, None), Ok(Some(_)) ) } @@ -168,6 +217,7 @@ pub(super) fn completed_data_usage_info( results: &[DataUsageCache], expected_sources: &HashSet, all_buckets: &[String], + tier_registry_names: &[String], bucket_plan_complete: bool, budget_elapsed: bool, cancelled: bool, @@ -183,6 +233,19 @@ pub(super) fn completed_data_usage_info( return None; } + // A generation is comparable across nodes because it is derived from the + // frozen registry names. Cycle and leader fencing remain separate cache + // metadata. Legacy peers omit the generation; an all-legacy result remains + // readable, but mixing legacy and new (or two new generations) would make + // the per-tier accounting ambiguous. + let registry_generation = results.first()?.info.tier_registry_generation; + if results.iter().any(|result| match registry_generation { + Some(generation) => result.info.tier_registry_generation != Some(generation), + None => result.info.tier_registry_generation.is_some(), + }) { + return None; + } + if results.iter().any(|result| result.root().is_none()) { return None; } @@ -200,9 +263,16 @@ pub(super) fn completed_data_usage_info( if !total.checked_merge(&merged) { return None; } + if !tier_accounting_proof_is_publishable(&merged, registry_generation, tier_registry_names) { + return None; + } buckets_usage.insert(bucket.clone(), checked_bucket_usage_info(&merged)?); } + if !tier_accounting_proof_is_publishable(&total, registry_generation, tier_registry_names) { + return None; + } + let merged_last_update = results.iter().filter_map(|result| result.info.last_update).max()?; let bucket_sizes = buckets_usage .iter() @@ -226,6 +296,109 @@ pub(super) fn completed_data_usage_info( Some((data_usage_info, merged_last_update)) } +fn tier_accounting_proof_is_publishable( + entry: &DataUsageEntry, + registry_generation: Option, + tier_registry_names: &[String], +) -> bool { + let has_scalar_usage = entry.size > 0 + || entry.objects > 0 + || entry.versions > 0 + || entry.delete_markers > 0 + || entry.failed_objects > 0 + || !entry.obj_sizes.is_empty() + || !entry.obj_versions.is_empty() + || entry.replication_stats.as_ref().is_some_and(|stats| !stats.is_empty()); + let has_tier_accounted_data = entry + .all_tier_stats + .as_ref() + .is_some_and(|stats| stats.tiers.values().any(|tier| !tier.is_empty())) + || entry.unknown_tier_stats.as_ref().is_some_and(|stats| { + stats.counter_overflowed + || stats.unknown_bytes > 0 + || stats.unknown_physical_bytes > 0 + || stats.unknown_objects > 0 + || stats.unknown_versions > 0 + }); + + let Some(proof) = entry.tier_accounting_proof else { + return !has_scalar_usage && !has_tier_accounted_data; + }; + if proof.overflowed + || entry + .unknown_tier_stats + .as_ref() + .is_some_and(|stats| stats.counter_overflowed) + || u64::try_from(entry.size).ok() != Some(proof.logical_total) + { + return false; + } + let unknown_logical = entry.unknown_tier_stats.as_ref().map_or(0, |stats| stats.unknown_bytes); + let unknown_physical = entry + .unknown_tier_stats + .as_ref() + .map_or(0, |stats| stats.unknown_physical_bytes); + if proof + .logical_known + .checked_add(unknown_logical) + .is_none_or(|total| total != proof.logical_total) + || proof + .physical_known + .checked_add(unknown_physical) + .is_none_or(|total| total != proof.physical_total) + { + return false; + } + if registry_generation.is_some() + && entry.all_tier_stats.as_ref().is_some_and(|stats| { + stats.tiers.keys().any(|tier| { + tier != crate::UNKNOWN_TIER + && tier != crate::storageclass::STANDARD + && tier != crate::storageclass::RRS + && !tier_registry_names.iter().any(|allowed| allowed == tier) + }) + }) + { + return false; + } + + if !has_tier_accounted_data { + return true; + } + + let Some(tiers) = entry.all_tier_stats.as_ref() else { + return false; + }; + let map_unknown_physical = tiers.tiers.get(crate::UNKNOWN_TIER).map_or(0, |stats| stats.total_size); + let companion_unknown_physical = entry + .unknown_tier_stats + .as_ref() + .map_or(0, |stats| stats.unknown_physical_bytes); + if map_unknown_physical != companion_unknown_physical { + return false; + } + + // A no-configuration scan intentionally stores only UNKNOWN_TIER after + // the first unknown object; STANDARD/RRS remain absent to preserve the + // historical empty-map shape. In that shape the scalar proof is the sole + // source of known physical bytes. Configured registries seed at least one + // non-UNKNOWN key, whose map total must match the proof. + let has_known_tier_map = tiers.tiers.keys().any(|tier| tier.as_str() != crate::UNKNOWN_TIER); + if !has_known_tier_map { + return true; + } + let Some(known_tier_physical_total) = tiers + .tiers + .iter() + .filter(|(tier, _)| tier.as_str() != crate::UNKNOWN_TIER) + .map(|(_, stats)| stats) + .try_fold(0_u64, |total, stats| total.checked_add(stats.total_size)) + else { + return false; + }; + proof.physical_known == known_tier_physical_total +} + pub(super) async fn send_cache_root_entry_info( bucket_result_tx: &mpsc::Sender, cache: &DataUsageCache, @@ -320,13 +493,14 @@ pub(super) async fn persist_and_publish_cache_snapshot( return None; } if matches!( - current_cache_root_entry( + current_cache_root_entry_with_generation( &persisted, DATA_USAGE_ROOT, source, cache_snapshot.info.next_cycle, cache_snapshot.info.leader_epoch, scan_plan_digest, + cache_snapshot.info.tier_registry_generation, ), Ok(Some(_)) ) { diff --git a/crates/scanner/src/scanner_io/io_cache.rs b/crates/scanner/src/scanner_io/io_cache.rs index aa0339e4b..5083d5783 100644 --- a/crates/scanner/src/scanner_io/io_cache.rs +++ b/crates/scanner/src/scanner_io/io_cache.rs @@ -31,6 +31,7 @@ impl ScannerIOCache for SetDisks { all_buckets, digest: scan_plan_digest, leader_epoch, + tier_registry_generation, dirty_usage_buckets, bucket_failures, pending_maintenance_work, @@ -48,6 +49,7 @@ impl ScannerIOCache for SetDisks { next_cycle: want_cycle, last_update: Some(now), leader_epoch, + tier_registry_generation: Some(tier_registry_generation), source: Some(source), snapshot_complete: true, scan_plan_digest: Some(scan_plan_digest), @@ -218,6 +220,13 @@ impl ScannerIOCache for SetDisks { "Scanner old data usage cache load failed; rebuilding from bucket caches" ); } + // Fence a stale set aggregate before copying entries into per-bucket work caches. + if old_cache.info.next_cycle <= want_cycle + && old_cache.info.leader_epoch <= leader_epoch + && old_cache.info.tier_registry_generation != Some(tier_registry_generation) + { + old_cache.info.scan_plan_digest = None; + } match old_cache.prepare_for_scan( DATA_USAGE_ROOT, want_cycle, @@ -267,6 +276,7 @@ impl ScannerIOCache for SetDisks { name: DATA_USAGE_ROOT.to_string(), next_cycle: want_cycle, leader_epoch, + tier_registry_generation: Some(tier_registry_generation), source: Some(source), snapshot_complete: false, scan_plan_digest: Some(scan_plan_digest), @@ -328,8 +338,9 @@ impl ScannerIOCache for SetDisks { }; let mut cache = cache_mutex_clone.lock().await; - apply_bucket_result_to_cache(&mut cache, result, SystemTime::now()); - completed_bucket_count_clone.fetch_add(1, Ordering::Relaxed); + if apply_bucket_result_to_cache(&mut cache, result, SystemTime::now()) { + completed_bucket_count_clone.fetch_add(1, Ordering::Relaxed); + } } } } @@ -460,6 +471,7 @@ impl ScannerIOCache for SetDisks { session_id: remote_session_id, session_sequence: request_sequence, scan_plan_digest: bucket_scan_plan_digest, + tier_registry_generation, skip_healing: healing, scan_mode, }, @@ -675,14 +687,17 @@ impl ScannerIOCache for SetDisks { continue; } }; - let scan_state = current_cache_root_or_prepare( + let scan_state = current_cache_root_or_prepare_with_generation( &mut cache, &bucket.name, source, want_cycle, leader_epoch, bucket_scan_plan_digest, - require_cache_source, + DataUsageCacheReuseOptions { + require_source: require_cache_source, + tier_registry_generation: Some(tier_registry_generation), + }, ); let outcome = match scan_state { DataUsageCacheScanState::Current(root) => { @@ -1108,6 +1123,7 @@ impl ScannerIOCache for SetDisks { name: DATA_USAGE_ROOT.to_string(), next_cycle: want_cycle, leader_epoch, + tier_registry_generation: Some(tier_registry_generation), source: Some(source), snapshot_complete: false, scan_plan_digest: Some(scan_plan_digest), diff --git a/crates/scanner/src/scanner_io/io_cycle.rs b/crates/scanner/src/scanner_io/io_cycle.rs index 64763655f..bf6590ecf 100644 --- a/crates/scanner/src/scanner_io/io_cycle.rs +++ b/crates/scanner/src/scanner_io/io_cycle.rs @@ -48,6 +48,7 @@ impl ScannerIOCycle for ECStore { scan_mode: HealScanMode, ) -> Result { let child_token = ctx.child_token(); + let _tier_cycle_guard = begin_tier_registry_cycle(want_cycle, leader_epoch); // Check the local pool metadata before listing buckets. A failed or // canceled decommission remains suspended after its worker exits, so @@ -127,6 +128,8 @@ impl ScannerIOCycle for ECStore { scanner_bucket_plan_digest(&all_buckets, crate::scanner::scanner_activity_snapshot_digest(&activity_before)); let dirty_usage_snapshot = Arc::new(snapshot_dirty_usage_buckets(&all_buckets, dirty_generation_before_bucket_list)); let cache_cycle_floor = Arc::new(AtomicU64::new(want_cycle)); + let tier_registry = runtime_tier_registry_for_cycle(want_cycle, leader_epoch).await; + let tier_registry_generation = tier_registry.generation; if all_buckets.is_empty() { reset_set_scan_gauges(); @@ -157,6 +160,9 @@ impl ScannerIOCycle for ECStore { { return Ok(ScannerCycleResult::new(status, None)); } + if status == ScannerCycleStatus::Complete { + complete_tier_registry_cycle(want_cycle, leader_epoch); + } let dirty_usage_clear = (status == ScannerCycleStatus::Complete).then(|| dirty_usage_snapshot.buckets.as_ref().clone()); let remote_dirty_usage_acknowledgements = if status == ScannerCycleStatus::Complete { @@ -249,6 +255,7 @@ impl ScannerIOCycle for ECStore { all_buckets: Arc::clone(&all_buckets), digest: scan_plan_digest, leader_epoch, + tier_registry_generation, dirty_usage_buckets: dirty_usage_snapshot.buckets.clone(), bucket_failures: bucket_failures.clone(), pending_maintenance_work: pending_maintenance_work.clone(), @@ -366,6 +373,7 @@ impl ScannerIOCycle for ECStore { &results, &expected_sources, &all_bucket_names, + &tier_registry.names, bucket_plan_complete, budget_elapsed, ctx.is_cancelled(), @@ -391,6 +399,9 @@ impl ScannerIOCycle for ECStore { &failed_buckets, ); result?; + if cycle_status == ScannerCycleStatus::Complete { + complete_tier_registry_cycle(want_cycle, leader_epoch); + } let remote_dirty_usage_acknowledgements = if cycle_status == ScannerCycleStatus::Complete { crate::scanner::scanner_dirty_usage_acknowledgements(&activity_before) } else { diff --git a/crates/scanner/src/scanner_io/io_disk.rs b/crates/scanner/src/scanner_io/io_disk.rs index 0887c6f7d..cf209206e 100644 --- a/crates/scanner/src/scanner_io/io_disk.rs +++ b/crates/scanner/src/scanner_io/io_disk.rs @@ -114,12 +114,13 @@ impl ScannerIODisk for Disk { .map(|v| ObjectInfo::from_file_info(v, item.bucket.as_str(), object_path.as_str(), versioned)) .collect::>(); - let mut size_summary = SizeSummary::default(); - // The caller supplies one registry snapshot for the whole folder scan; // seeding from it prevents a TTL refresh from mixing generations in a // single result. - size_summary.tier_stats = tier_stats_template(tier_names); + let mut size_summary = SizeSummary { + tier_stats: tier_stats_template(tier_names), + ..Default::default() + }; let lock_config = object_lock_config_for_scanner_item(&item).await; @@ -134,7 +135,9 @@ impl ScannerIODisk for Disk { if !free_version_infos.is_empty() { for oi in free_version_infos { - enqueue_runtime_free_version(oi).await; + if ScannerItem::tier_is_known(&oi, tier_names) { + enqueue_runtime_free_version(oi).await; + } } } diff --git a/crates/scanner/src/scanner_io/publish_gate_tests.rs b/crates/scanner/src/scanner_io/publish_gate_tests.rs index c6acdea1a..a4a5f3ea4 100644 --- a/crates/scanner/src/scanner_io/publish_gate_tests.rs +++ b/crates/scanner/src/scanner_io/publish_gate_tests.rs @@ -13,7 +13,8 @@ // limitations under the License. use super::*; -use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage}; +use crate::data_usage_define::{UNKNOWN_TIER, UnknownTierStats, hash_path}; +use rustfs_data_usage::{ReplicationAllStats, ReplicationTargetUsage, TierAccountingProof}; const TEST_PLAN_DIGEST: DataUsageScanPlanDigest = DataUsageScanPlanDigest([7; 32]); @@ -89,6 +90,11 @@ fn completed_root_cache(bucket: &str, objects: usize, update_secs: u64, source: DataUsageEntry { objects, size: objects.saturating_mul(10), + tier_accounting_proof: Some(TierAccountingProof { + logical_total: u64::try_from(objects.saturating_mul(10)).unwrap_or(u64::MAX), + logical_known: u64::try_from(objects.saturating_mul(10)).unwrap_or(u64::MAX), + ..Default::default() + }), ..Default::default() }, ); @@ -102,7 +108,7 @@ fn completed_data_usage_info_for_test( cancelled: bool, ) -> Option<(DataUsageInfo, SystemTime)> { let expected_sources = results.iter().filter_map(|result| result.info.source).collect::>(); - completed_data_usage_info(results, &expected_sources, all_buckets, true, budget_elapsed, cancelled) + completed_data_usage_info(results, &expected_sources, all_buckets, &[], true, budget_elapsed, cancelled) } #[test] @@ -122,11 +128,21 @@ fn completed_data_usage_info_publishes_tier_stats_across_sets() { let mut first_set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); let mut tiered = DataUsageEntry::default(); tiered.add_tier_sizes(&warm(100, 2, 1)); + tiered.tier_accounting_proof = Some(TierAccountingProof { + physical_total: 100, + physical_known: 100, + ..Default::default() + }); first_set.replace("bucket-b", DATA_USAGE_ROOT, tiered); let mut second_set = completed_root_cache("bucket-b", 2, 20, DataUsageCacheSource::new(1, 0)); let mut tiered = DataUsageEntry::default(); tiered.add_tier_sizes(&warm(50, 1, 1)); + tiered.tier_accounting_proof = Some(TierAccountingProof { + physical_total: 50, + physical_known: 50, + ..Default::default() + }); second_set.replace("bucket-a", DATA_USAGE_ROOT, tiered); let (data_usage_info, _) = completed_data_usage_info_for_test(&[first_set, second_set], &all_buckets, false, false) @@ -145,6 +161,269 @@ fn completed_data_usage_info_publishes_tier_stats_across_sets() { ); } +#[test] +fn completed_data_usage_info_rejects_logical_proof_mismatch() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.add_tier_sizes(&HashMap::from([( + "WARM".to_string(), + TierStats { + total_size: 10, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 10, + logical_known: 9, + physical_total: 10, + physical_known: 10, + ..Default::default() + }); + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + +#[test] +fn completed_data_usage_info_rejects_logical_total_size_mismatch() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.size = 11; + entry.add_tier_sizes(&HashMap::from([( + "WARM".to_string(), + TierStats { + total_size: 10, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 10, + logical_known: 10, + physical_total: 10, + physical_known: 10, + ..Default::default() + }); + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + +#[test] +fn completed_data_usage_info_rejects_physical_proof_mismatch() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.add_tier_sizes(&HashMap::from([( + "WARM".to_string(), + TierStats { + total_size: 10, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 10, + logical_known: 10, + physical_total: 9, + physical_known: 9, + ..Default::default() + }); + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + +#[test] +fn completed_data_usage_info_rejects_unknown_physical_double_accounting() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.add_tier_sizes(&HashMap::from([( + UNKNOWN_TIER.to_string(), + TierStats { + total_size: 10, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.add_unknown_tier_stats(&UnknownTierStats { + unknown_physical_bytes: 9, + ..Default::default() + }); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 10, + logical_known: 10, + physical_total: 10, + physical_known: 10, + ..Default::default() + }); + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + +#[test] +fn completed_data_usage_info_rejects_unknown_counter_overflow() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.add_unknown_tier_stats(&UnknownTierStats { + counter_overflowed: true, + ..Default::default() + }); + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + +#[test] +fn completed_data_usage_info_accepts_no_tier_standard_empty_map_with_proof() { + let all_buckets = vec!["bucket-a".to_string()]; + let set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + + let (info, _) = completed_data_usage_info_for_test(&[set], &all_buckets, false, false) + .expect("no-tier STANDARD/RRS usage does not require a tier map"); + assert!(info.tier_stats.is_none()); +} + +#[test] +fn completed_data_usage_info_accepts_no_tier_unknown_and_standard_shape() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.size = 13; + entry.add_tier_sizes(&HashMap::from([( + UNKNOWN_TIER.to_string(), + TierStats { + total_size: 3, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.add_unknown_tier_stats(&UnknownTierStats { + unknown_bytes: 9, + unknown_physical_bytes: 3, + unknown_objects: 1, + unknown_versions: 1, + ..Default::default() + }); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 13, + logical_known: 4, + physical_total: 7, + physical_known: 4, + ..Default::default() + }); + + let (info, _) = completed_data_usage_info_for_test(&[set], &all_buckets, false, false) + .expect("no-tier STANDARD plus UNKNOWN should remain publishable"); + assert_eq!( + info.tier_stats.expect("unknown bucket should be retained").tiers[UNKNOWN_TIER].total_size, + 3 + ); +} + +#[test] +fn completed_data_usage_info_accepts_unknown_only_with_current_registry_generation() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + set.info.tier_registry_generation = Some(7); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.size = 13; + entry.add_tier_sizes(&HashMap::from([( + UNKNOWN_TIER.to_string(), + TierStats { + total_size: 3, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.add_unknown_tier_stats(&UnknownTierStats { + unknown_bytes: 9, + unknown_physical_bytes: 3, + unknown_objects: 1, + unknown_versions: 1, + ..Default::default() + }); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 13, + logical_known: 4, + physical_total: 7, + physical_known: 4, + ..Default::default() + }); + + let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]); + assert!( + completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_some() + ); +} + +#[test] +fn completed_data_usage_info_rejects_non_registry_tier_in_current_generation() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + set.info.tier_registry_generation = Some(7); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.add_tier_sizes(&HashMap::from([ + ( + "WARM".to_string(), + TierStats { + total_size: 4, + num_versions: 1, + num_objects: 1, + }, + ), + ( + "RETIRED".to_string(), + TierStats { + total_size: 6, + num_versions: 1, + num_objects: 1, + }, + ), + ])); + entry.tier_accounting_proof = Some(TierAccountingProof { + logical_total: 10, + logical_known: 10, + physical_total: 10, + physical_known: 10, + ..Default::default() + }); + + let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]); + assert!( + completed_data_usage_info(&[set], &expected_sources, &all_buckets, &["WARM".to_string()], true, false, false,).is_none() + ); +} + +#[test] +fn completed_data_usage_info_rejects_legacy_proof_missing_when_tier_accounted() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.add_tier_sizes(&HashMap::from([( + "WARM".to_string(), + TierStats { + total_size: 10, + num_versions: 1, + num_objects: 1, + }, + )])); + entry.tier_accounting_proof = None; + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + +#[test] +fn completed_data_usage_info_rejects_legacy_proof_missing_for_scalar_usage() { + let all_buckets = vec!["bucket-a".to_string()]; + let mut set = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let entry = set.cache.get_mut(&hash_path("bucket-a").key()).expect("bucket entry"); + entry.tier_accounting_proof = None; + + assert!(completed_data_usage_info_for_test(&[set], &all_buckets, false, false).is_none()); +} + #[test] fn completed_data_usage_info_omits_tier_stats_without_tiered_objects() { let all_buckets = vec!["bucket-a".to_string()]; @@ -156,6 +435,51 @@ fn completed_data_usage_info_omits_tier_stats_without_tiered_objects() { assert!(data_usage_info.tier_stats.is_none()); } +#[test] +fn completed_data_usage_info_rejects_legacy_and_new_tier_generations_mixed() { + let all_buckets = vec!["bucket-a".to_string()]; + let legacy = completed_root_cache("bucket-a", 1, 10, DataUsageCacheSource::new(0, 0)); + let mut current = completed_root_cache("bucket-a", 1, 20, DataUsageCacheSource::new(1, 0)); + current.info.tier_registry_generation = Some(42); + + assert!( + completed_data_usage_info_for_test(&[legacy, current], &all_buckets, false, false).is_none(), + "legacy and generation-tagged sets must not publish a mixed snapshot" + ); +} + +#[test] +fn current_cache_root_with_new_tier_generation_resets_old_cache() { + let source = DataUsageCacheSource::new(0, 0); + let mut cache = completed_root_cache("bucket-a", 1, 10, source); + cache.info.tier_registry_generation = Some(1); + + let state = current_cache_root_or_prepare_with_generation( + &mut cache, + DATA_USAGE_ROOT, + source, + 0, + 0, + TEST_PLAN_DIGEST, + DataUsageCacheReuseOptions { + require_source: false, + tier_registry_generation: Some(2), + }, + ); + + assert!(matches!( + state, + DataUsageCacheScanState::Prepared { + outcome: DataUsageCachePrepareOutcome::Reset, + .. + } + )); + assert!(cache.cache.is_empty(), "old-generation entries must not be reused"); + assert_eq!(cache.info.tier_registry_generation, None); + assert_eq!(cache.info.scan_plan_digest, Some(TEST_PLAN_DIGEST)); + assert!(!cache.info.snapshot_complete); +} + #[test] fn completed_data_usage_info_requires_every_set_before_publish() { let all_buckets = vec!["bucket-a".to_string(), "bucket-b".to_string(), "bucket-empty".to_string()]; @@ -280,6 +604,13 @@ fn completed_data_usage_info_flattens_nested_bucket_entries() { replica_size: 2048, replica_count: 2, }), + tier_accounting_proof: Some(TierAccountingProof { + logical_total: 2048, + logical_known: 2048, + physical_total: 2048, + physical_known: 2048, + ..Default::default() + }), ..Default::default() }; nested.obj_sizes.add(2048); @@ -347,7 +678,8 @@ fn completed_data_usage_info_requires_exact_topology_sources() { let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0), DataUsageCacheSource::new(1, 0)]); assert!( - completed_data_usage_info(&[first_set, unexpected_set], &expected_sources, &all_buckets, true, false, false).is_none() + completed_data_usage_info(&[first_set, unexpected_set], &expected_sources, &all_buckets, &[], true, false, false) + .is_none() ); } @@ -357,7 +689,7 @@ fn completed_data_usage_info_rejects_incomplete_bucket_plan() { let set = completed_root_cache("bucket", 2, 10, DataUsageCacheSource::new(0, 0)); let expected_sources = HashSet::from([DataUsageCacheSource::new(0, 0)]); - assert!(completed_data_usage_info(&[set], &expected_sources, &all_buckets, false, false, false).is_none()); + assert!(completed_data_usage_info(&[set], &expected_sources, &all_buckets, &[], false, false, false).is_none()); } #[test] diff --git a/crates/scanner/src/scanner_io/tests.rs b/crates/scanner/src/scanner_io/tests.rs index 9ed077b41..d52d360a0 100644 --- a/crates/scanner/src/scanner_io/tests.rs +++ b/crates/scanner/src/scanner_io/tests.rs @@ -1328,7 +1328,7 @@ fn apply_bucket_result_to_cache_updates_bucket_entry() { ); let update_time = SystemTime::now(); - apply_bucket_result_to_cache( + assert!(apply_bucket_result_to_cache( &mut cache, DataUsageEntryInfo { name: "bucket".to_string(), @@ -1338,12 +1338,51 @@ fn apply_bucket_result_to_cache_updates_bucket_entry() { objects: 2, ..Default::default() }, + tier_registry_generation: None, }, update_time, - ); + )); assert_eq!(cache.info.last_update, Some(update_time)); let entry = cache.find("bucket").expect("bucket entry should remain present"); assert_eq!(entry.size, 10); assert_eq!(entry.objects, 2); } + +#[test] +fn apply_bucket_result_to_cache_rejects_a_different_tier_generation() { + let mut cache = DataUsageCache { + info: DataUsageCacheInfo { + name: DATA_USAGE_ROOT.to_string(), + tier_registry_generation: Some(7), + ..Default::default() + }, + ..Default::default() + }; + cache.replace( + "bucket", + DATA_USAGE_ROOT, + DataUsageEntry { + size: 3, + ..Default::default() + }, + ); + + let applied = apply_bucket_result_to_cache( + &mut cache, + DataUsageEntryInfo { + name: "bucket".to_string(), + parent: DATA_USAGE_ROOT.to_string(), + entry: DataUsageEntry { + size: 11, + ..Default::default() + }, + tier_registry_generation: Some(8), + }, + SystemTime::now(), + ); + + assert!(!applied); + assert_eq!(cache.find("bucket").map(|entry| entry.size), Some(3)); + assert!(cache.info.last_update.is_none()); +} diff --git a/crates/storage-api/src/lib.rs b/crates/storage-api/src/lib.rs index 1114349e1..7768dad62 100644 --- a/crates/storage-api/src/lib.rs +++ b/crates/storage-api/src/lib.rs @@ -46,6 +46,7 @@ pub const NS_SCANNER_SERVER_EPOCH_QUERY: &str = "ns_scanner_server_epoch"; pub const NS_SCANNER_SESSION_ID_QUERY: &str = "ns_scanner_session_id"; pub const NS_SCANNER_SESSION_SEQUENCE_QUERY: &str = "ns_scanner_session_sequence"; pub const NS_SCANNER_PROTOCOL_VERSION_QUERY: &str = "ns_scanner_protocol"; +pub const NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY: &str = "ns_scanner_tier_registry_generation"; pub const NS_SCANNER_PROTOCOL_VERSION: u16 = 3; pub const SCANNER_ACTIVITY_LEGACY_PROTOCOL_VERSION: u32 = 0; pub const SCANNER_ACTIVITY_PREVIOUS_PROTOCOL_VERSION: u32 = 5; @@ -57,6 +58,8 @@ pub struct NsScannerCapabilityResponse { pub version: u16, pub server_epoch: uuid::Uuid, pub proof: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub supports_tier_registry_generation: Option, } pub mod admin; diff --git a/rustfs/src/storage/rpc/http_service.rs b/rustfs/src/storage/rpc/http_service.rs index c2983cc43..052dbaff5 100644 --- a/rustfs/src/storage/rpc/http_service.rs +++ b/rustfs/src/storage/rpc/http_service.rs @@ -19,13 +19,15 @@ use crate::storage::storage_api::rpc_consumer::http_service::{ DEFAULT_READ_BUFFER_SIZE, DeleteOptions, DiskStore, NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, PUT_FILE_AUTH_TRAILER_LEN, PUT_FILE_AUTH_V1, PUT_FILE_CAPABILITY_VERSION, PutFileCapabilityResponse, StorageDiskRpcExt as _, WALK_DIR_STREAM_COMPLETION_V1, WalkDirOptions, check_and_record_signed_rpc_nonce, find_local_disk_by_ref, - sign_ns_scanner_capability, sign_put_file_capability, verify_put_file_auth_trailer, verify_rpc_signature, + sign_ns_scanner_capability_with_tier_registry_generation, sign_put_file_capability, verify_put_file_auth_trailer, + verify_rpc_signature, }; #[cfg(test)] use crate::storage::storage_api::rpc_consumer::http_service::{ NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY, NS_SCANNER_LEADER_EPOCH_QUERY, NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY, NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, - PUT_FILE_CAPABILITY_CHALLENGE_QUERY, PUT_FILE_CAPABILITY_QUERY, WALK_DIR_BODY_SHA256_QUERY, + NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY, PUT_FILE_CAPABILITY_CHALLENGE_QUERY, PUT_FILE_CAPABILITY_QUERY, + WALK_DIR_BODY_SHA256_QUERY, }; use crate::storage::storage_api::runtime_sources_consumer::runtime_sources; use crate::storage::storage_api::tonic_rpc_auth_failure_reason; @@ -290,6 +292,8 @@ struct NsScannerQuery { struct NsScannerCapabilityQuery { ns_scanner_protocol: Option, ns_scanner_challenge: Option, + #[serde(rename = "ns_scanner_tier_registry_generation")] + ns_scanner_tier_registry_generation: Option, } fn verify_ns_scanner_body_digest(query: &NsScannerQuery, body: &[u8]) -> bool { @@ -410,7 +414,9 @@ async fn handle_internode_rpc(req: Request) -> Response { (Method::GET, WALK_DIR_PATH) | (Method::HEAD, WALK_DIR_PATH) => handle_walk_dir(req).await, (Method::GET, NS_SCANNER_PATH) => match parse_query::(&req) { Ok(query) if query.ns_scanner_protocol == Some(NS_SCANNER_PROTOCOL_VERSION) => match query.ns_scanner_challenge { - Some(challenge) if !challenge.is_nil() => ns_scanner_capability_response(challenge), + Some(challenge) if !challenge.is_nil() => { + ns_scanner_capability_response(challenge, query.ns_scanner_tier_registry_generation == Some(true)) + } Some(_) | None => response_with_status(StatusCode::BAD_REQUEST, "namespace scanner challenge is invalid"), }, Ok(_) => response_with_status(StatusCode::UPGRADE_REQUIRED, "namespace scanner protocol is unsupported"), @@ -466,31 +472,34 @@ fn record_internode_rpc_error(operation: Option<&'static str>) { } } -fn ns_scanner_capability_response(challenge: uuid::Uuid) -> Response { +fn ns_scanner_capability_response(challenge: uuid::Uuid, include_tier_registry_generation: bool) -> Response { let server_epoch = *NS_SCANNER_SERVER_EPOCH; - let proof = match sign_ns_scanner_capability(challenge, server_epoch) { - Ok(proof) => proof, - Err(err) => { - error!( - event = EVENT_RPC_REQUEST_FAILED, - component = LOG_COMPONENT_INTERNODE_RPC, - subsystem = LOG_SUBSYSTEM_NAMESPACE_SCANNER, - operation = INTERNODE_OPERATION_NS_SCANNER, - result = "failed", - status_code = StatusCode::UPGRADE_REQUIRED.as_u16(), - rpc_path = NS_SCANNER_PATH, - method = %Method::GET, - reason = "capability_authentication_unavailable", - error = %err, - "internode rpc request failed" - ); - return response_with_status(StatusCode::UPGRADE_REQUIRED, "namespace scanner RPC authentication is unavailable"); - } - }; + let proof = + match sign_ns_scanner_capability_with_tier_registry_generation(challenge, server_epoch, include_tier_registry_generation) + { + Ok(proof) => proof, + Err(err) => { + error!( + event = EVENT_RPC_REQUEST_FAILED, + component = LOG_COMPONENT_INTERNODE_RPC, + subsystem = LOG_SUBSYSTEM_NAMESPACE_SCANNER, + operation = INTERNODE_OPERATION_NS_SCANNER, + result = "failed", + status_code = StatusCode::UPGRADE_REQUIRED.as_u16(), + rpc_path = NS_SCANNER_PATH, + method = %Method::GET, + reason = "capability_authentication_unavailable", + error = %err, + "internode rpc request failed" + ); + return response_with_status(StatusCode::UPGRADE_REQUIRED, "namespace scanner RPC authentication is unavailable"); + } + }; let body = match rmp_serde::to_vec_named(&NsScannerCapabilityResponse { version: NS_SCANNER_PROTOCOL_VERSION, server_epoch, proof, + supports_tier_registry_generation: include_tier_registry_generation.then_some(true), }) { Ok(body) => body, Err(err) => { @@ -1676,12 +1685,13 @@ mod tests { LOG_SUBSYSTEM_NAMESPACE_SCANNER, LOG_SUBSYSTEM_ROUTING, NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY, NS_SCANNER_LEADER_EPOCH_QUERY, NS_SCANNER_PATH, NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY, NS_SCANNER_SESSION_ID_QUERY, - NS_SCANNER_SESSION_SEQUENCE_QUERY, NsScannerQuery, PUT_FILE_AUTH_STREAM_PATH, PUT_FILE_CAPABILITY_PATH, - PUT_FILE_STREAM_PATH, PutFileQuery, READ_FILE_STREAM_PATH, WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_PATH, WalkDirQuery, - append_walk_dir_completion, internode_http_operation, internode_rpc_subsystem, is_internode_rpc_path, - ns_scanner_response_body, ns_scanner_server_epoch_matches, put_body_size_mismatch, put_file_auth_nonce, - put_file_capability_response, put_file_server_epoch_matches, put_file_stage_error_message, put_file_target_lock, - read_file_body_stream, read_file_stream_buffer_size, remote_scanner_claim_rejection, response_with_disk_error, + NS_SCANNER_SESSION_SEQUENCE_QUERY, NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY, NsScannerCapabilityResponse, + NsScannerQuery, PUT_FILE_AUTH_STREAM_PATH, PUT_FILE_CAPABILITY_PATH, PUT_FILE_STREAM_PATH, PutFileQuery, + READ_FILE_STREAM_PATH, WALK_DIR_BODY_SHA256_QUERY, WALK_DIR_PATH, WalkDirQuery, append_walk_dir_completion, + internode_http_operation, internode_rpc_subsystem, is_internode_rpc_path, ns_scanner_response_body, + ns_scanner_server_epoch_matches, put_body_size_mismatch, put_file_auth_nonce, put_file_capability_response, + put_file_server_epoch_matches, put_file_stage_error_message, put_file_target_lock, read_file_body_stream, + read_file_stream_buffer_size, remote_scanner_claim_rejection, response_with_disk_error, supports_walk_dir_stream_completion, validate_walk_dir_completion_request, verify_internode_rpc_signature, verify_ns_scanner_body_digest, verify_walk_dir_body_digest, walk_dir_response_body, write_authenticated_put_file, write_body_chunks_to_writer, write_put_file_body_chunks_to_writer, @@ -2114,6 +2124,51 @@ mod tests { ); assert!(serde_urlencoded::from_str::(&query).is_err()); assert!(serde_urlencoded::from_str::("ns_scanner_protocol=1&unexpected=true").is_err()); + let marked = + format!("ns_scanner_protocol=3&ns_scanner_challenge={request_id}&{NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY}=true"); + assert_eq!( + serde_urlencoded::from_str::(&marked) + .expect("generation marker should be accepted") + .ns_scanner_tier_registry_generation, + Some(true) + ); + } + + #[test] + fn namespace_scanner_capability_response_support_is_optional_for_old_peers() { + #[derive(serde::Deserialize)] + #[serde(deny_unknown_fields)] + struct LegacyCapabilityResponse { + version: u16, + server_epoch: uuid::Uuid, + proof: Vec, + } + + let old = rmp_serde::to_vec_named(&NsScannerCapabilityResponse { + version: super::NS_SCANNER_PROTOCOL_VERSION, + server_epoch: uuid::Uuid::new_v4(), + proof: vec![1, 2, 3], + supports_tier_registry_generation: None, + }) + .expect("old response shape should encode"); + let decoded: NsScannerCapabilityResponse = rmp_serde::from_slice(&old).expect("old response should decode"); + assert_eq!(decoded.supports_tier_registry_generation, None); + let legacy_decoded: LegacyCapabilityResponse = + rmp_serde::from_slice(&old).expect("legacy reader should decode old shape"); + assert_eq!(legacy_decoded.version, super::NS_SCANNER_PROTOCOL_VERSION); + assert!(!legacy_decoded.server_epoch.is_nil()); + assert_eq!(legacy_decoded.proof, vec![1, 2, 3]); + + let current = rmp_serde::to_vec_named(&NsScannerCapabilityResponse { + version: super::NS_SCANNER_PROTOCOL_VERSION, + server_epoch: uuid::Uuid::new_v4(), + proof: vec![1, 2, 3], + supports_tier_registry_generation: Some(true), + }) + .expect("new response shape should encode"); + let decoded: NsScannerCapabilityResponse = rmp_serde::from_slice(¤t).expect("new response should decode"); + assert_eq!(decoded.supports_tier_registry_generation, Some(true)); + assert!(rmp_serde::from_slice::(¤t).is_err()); } #[test] diff --git a/rustfs/src/storage/storage_api.rs b/rustfs/src/storage/storage_api.rs index e5277c0dd..3a7eb6121 100644 --- a/rustfs/src/storage/storage_api.rs +++ b/rustfs/src/storage/storage_api.rs @@ -224,8 +224,8 @@ pub(crate) mod rpc_consumer { pub(crate) use super::super::storage_contracts::{ NS_SCANNER_BODY_SHA256_QUERY, NS_SCANNER_CAPABILITY_CHALLENGE_QUERY, NS_SCANNER_CYCLE_QUERY, NS_SCANNER_LEADER_EPOCH_QUERY, NS_SCANNER_REQUEST_ID_QUERY, NS_SCANNER_SERVER_EPOCH_QUERY, - NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, PUT_FILE_CAPABILITY_CHALLENGE_QUERY, - PUT_FILE_CAPABILITY_QUERY, WALK_DIR_BODY_SHA256_QUERY, + NS_SCANNER_SESSION_ID_QUERY, NS_SCANNER_SESSION_SEQUENCE_QUERY, NS_SCANNER_TIER_REGISTRY_GENERATION_QUERY, + PUT_FILE_CAPABILITY_CHALLENGE_QUERY, PUT_FILE_CAPABILITY_QUERY, WALK_DIR_BODY_SHA256_QUERY, }; pub(crate) use super::super::storage_contracts::{ NS_SCANNER_PROTOCOL_VERSION, NsScannerCapabilityResponse, PUT_FILE_AUTH_TRAILER_LEN, PUT_FILE_AUTH_V1, @@ -233,8 +233,8 @@ pub(crate) mod rpc_consumer { }; pub(crate) use super::super::{ DeleteOptions, DiskStore, StorageDiskRpcExt, WalkDirOptions, check_and_record_signed_rpc_nonce, - find_local_disk_by_ref, sign_ns_scanner_capability, sign_put_file_capability, verify_put_file_auth_trailer, - verify_rpc_signature, + find_local_disk_by_ref, sign_ns_scanner_capability_with_tier_registry_generation, sign_put_file_capability, + verify_put_file_auth_trailer, verify_rpc_signature, }; } @@ -520,9 +520,10 @@ pub(crate) mod ecstore_rpc { pub(crate) use rustfs_ecstore::api::rpc::{ KMS_SIGNAL_SUBSYSTEM, LocalPeerS3Client, PEER_RESTDRY_RUN, PEER_RESTSIGNAL, PEER_RESTSUB_SYS, PeerRestClient, PeerS3Client, SERVICE_SIGNAL_REFRESH_CONFIG, SERVICE_SIGNAL_RELOAD_DYNAMIC, TONIC_RPC_PREFIX, - check_and_record_signed_rpc_nonce, normalize_tonic_rpc_audience, sign_ns_scanner_capability, sign_put_file_capability, - sign_tonic_rpc_response_proof, tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers, - tonic_rpc_auth_failure_reason, verify_put_file_auth_trailer, verify_rpc_signature, verify_tonic_canonical_body_digest, + check_and_record_signed_rpc_nonce, normalize_tonic_rpc_audience, + sign_ns_scanner_capability_with_tier_registry_generation, sign_put_file_capability, sign_tonic_rpc_response_proof, + tonic_boot_epoch_challenge, tonic_boot_epoch_response_headers, tonic_rpc_auth_failure_reason, + verify_put_file_auth_trailer, verify_rpc_signature, verify_tonic_canonical_body_digest, verify_tonic_mutation_body_digest, verify_tonic_rpc_signature_with_bootstrap, }; #[cfg(test)] @@ -1727,8 +1728,16 @@ pub(crate) fn verify_put_file_auth_trailer( ecstore_rpc::verify_put_file_auth_trailer(url, method, nonce, trailer) } -pub(crate) fn sign_ns_scanner_capability(challenge: uuid::Uuid, server_epoch: uuid::Uuid) -> std::io::Result> { - ecstore_rpc::sign_ns_scanner_capability(challenge, server_epoch) +pub(crate) fn sign_ns_scanner_capability_with_tier_registry_generation( + challenge: uuid::Uuid, + server_epoch: uuid::Uuid, + supports_tier_registry_generation: bool, +) -> std::io::Result> { + ecstore_rpc::sign_ns_scanner_capability_with_tier_registry_generation( + challenge, + server_epoch, + supports_tier_registry_generation, + ) } pub(crate) fn sign_put_file_capability(