Compare commits

..

5 Commits

Author SHA1 Message Date
houseme 609c4b51b0 Merge branch 'main' into cxymds/fix-1936-unknown-tier
Signed-off-by: houseme <housemecn@gmail.com>
2026-08-23 12:33:38 +08:00
马登山 b0fe2b277a merge main and fix lint 2026-08-23 07:46:52 +08:00
马登山 fe3e779977 fix(scanner): fence unknown tier accounting 2026-08-23 05:45:30 +08:00
马登山 2ab23980f9 docs(scanner): clarify unknown tier stats 2026-08-22 23:42:58 +08:00
马登山 0b0c7ca4d7 fix(scanner): bound retired tier accounting 2026-08-22 23:41:31 +08:00
31 changed files with 2698 additions and 516 deletions
-20
View File
@@ -1,20 +0,0 @@
# Report-only calibration baseline from https://github.com/rustfs/rustfs/actions/runs/29394996173.
# Update counts only with a linked coverage run and a reviewed explanation.
phase = "report-only"
allowed_drop_percentage_points = 1.0
[crates."crates/iam"]
covered = 5149
count = 8131
[crates."crates/kms"]
covered = 2950
count = 4200
[crates."crates/policy"]
covered = 4636
count = 5464
[crates."crates/crypto"]
covered = 469
count = 494
-1
View File
@@ -36,7 +36,6 @@ script-tests: ## Run shell script tests
./scripts/test_manual_transition_runbooks.sh
./scripts/check_embedded_secrets.sh --self-test
python3 ./scripts/check_test_wiring.py --self-test
python3 ./scripts/check_security_coverage.py --self-test
python3 ./scripts/check_scheduled_validation_freshness.py --self-test
python3 ./scripts/s3-tests/test_report_compat.py
bash -n ./scripts/validate_object_data_cache_cold_stampede.sh
+12 -28
View File
@@ -12,12 +12,14 @@
# See the License for the specific language governing permissions and
# limitations under the License.
# Workspace line-coverage baseline and security-crate calibration
# (backlog#1153 infra-5/infra-6).
# Weekly workspace line-coverage baseline (backlog#1153 infra-5).
#
# NON-BLOCKING by design: the weekly job gives coverage a visible baseline and
# trend, while relevant pull requests run a report-only security-crate
# comparison. Neither job is a required check during calibration.
# NON-BLOCKING by design: this workflow only runs on schedule and manual
# dispatch, so it never attaches a status to a PR and must never be made a
# required check. It exists to give coverage a visible baseline and trend
# (per-crate table in the job summary, lcov artifact kept 90 days) — the
# per-crate ratchet for the security-critical crates builds on it later
# (backlog#1153 infra-6, report-only first per the ci-11 ladder).
#
# Measurement scope matches the PR test gate (ci.yml "Run tests"):
# `--workspace --exclude e2e_test` with the `ci` nextest profile. Doctests are
@@ -29,17 +31,6 @@
name: coverage
on:
pull_request:
branches: [main]
paths:
- "crates/iam/**"
- "crates/kms/**"
- "crates/policy/**"
- "crates/crypto/**"
- ".config/coverage-baselines.toml"
- "scripts/coverage_per_crate.py"
- "scripts/check_security_coverage.py"
- ".github/workflows/coverage.yml"
workflow_dispatch:
schedule:
# 07:00 UTC Sunday — staggered clear of the other Sunday crons: ci (00:00),
@@ -48,10 +39,6 @@ on:
# e2e-replication-nightly (04:00) and performance-ab (06:00) lanes.
- cron: "43 7 * * 0"
concurrency:
group: ${{ github.workflow }}-${{ github.event_name }}-${{ github.event.pull_request.number || github.ref }}
cancel-in-progress: ${{ github.event_name != 'schedule' }}
# Only alert-on-failure needs more than read access; it declares its own
# job-level `issues: write`.
permissions:
@@ -59,13 +46,12 @@ permissions:
jobs:
coverage:
name: Workspace line coverage
name: Workspace coverage (weekly)
runs-on: sm-standard-4
# The instrumented build cannot reuse the regular CI cache (different
# RUSTFLAGS), so a cold run rebuilds the workspace before running the
# full suite. Exact-head run 32573798257 needed 119m42s including reports
# and artifact upload, so keep a bounded 30-minute publication margin.
timeout-minutes: 150
# RUSTFLAGS), so a cold week rebuilds the workspace before running the
# full suite; give it double the test job's 60-minute budget.
timeout-minutes: 120
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
# Match the PR gate's nextest semantics (ci.yml runs `--profile ci`):
@@ -105,9 +91,7 @@ jobs:
cargo llvm-cov report --json --output-path target/llvm-cov/coverage.json
- name: Write per-crate summary
run: |
python3 scripts/coverage_per_crate.py target/llvm-cov/coverage.json >> "$GITHUB_STEP_SUMMARY"
python3 scripts/check_security_coverage.py target/llvm-cov/coverage.json >> "$GITHUB_STEP_SUMMARY"
run: python3 scripts/coverage_per_crate.py target/llvm-cov/coverage.json >> "$GITHUB_STEP_SUMMARY"
- name: Upload coverage artifact
if: always()
+676 -3
View File
@@ -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},
@@ -48,6 +53,11 @@ pub const DATA_USAGE_OBSERVED_OBJECT_NAME: &str = ".usage.observed.json";
// RUSTFS_COMPAT_TODO(scanner-usage-v2): keep .usage.json readable and removable during rolling upgrades from pre-v2 scanners. Remove after supported direct-upgrade sources all write .usage.v2.json.
pub const LEGACY_DATA_USAGE_OBJECT_NAME: &str = ".usage.json";
/// Fixed bucket for objects whose storage class is not in the scanner's
/// cycle-local tier registry. Keeping this key fixed prevents untrusted or
/// stale tier names from growing persisted per-tier maps without bound.
pub const UNKNOWN_TIER: &str = "UNKNOWN_TIER";
/// Returns true when `existing_last_update` is ahead of `now` by more than
/// [`USAGE_LAST_UPDATE_FUTURE_TOLERANCE`], i.e. the persisted timestamp cannot be
/// trusted for staleness comparisons and a fresh snapshot save must be allowed.
@@ -78,12 +88,301 @@ impl TierStats {
&& self.num_objects.checked_add(u.num_objects).is_some()
}
/// Add tier counters without allowing a counter to wrap.
pub fn checked_add(&self, u: &TierStats) -> Option<TierStats> {
Some(TierStats {
total_size: self.total_size.checked_add(u.total_size)?,
num_versions: self.num_versions.checked_add(u.num_versions)?,
num_objects: self.num_objects.checked_add(u.num_objects)?,
})
}
/// True when this tier contributed nothing, i.e. merging it is a no-op.
pub fn is_empty(&self) -> bool {
self.total_size == 0 && self.num_versions == 0 && self.num_objects == 0
}
}
/// Bounded diagnostics for objects whose tier is absent from the cycle
/// registry. Counters are authoritative; diagnostics are only a small,
/// redacted reconciliation aid and may be dropped at the configured caps.
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, PartialEq, Eq)]
pub struct UnknownTierStats {
/// Logical bytes retained in the scanner's normal usage total.
pub unknown_bytes: u64,
/// Physical bytes recorded in the per-tier accounting dimension.
///
/// Older writers only had `unknown_bytes`; decoding those snapshots keeps
/// this field at zero and the scanner fills it for new observations.
#[serde(default)]
pub unknown_physical_bytes: u64,
pub unknown_objects: u64,
pub unknown_versions: u64,
pub diagnostics_dropped: u64,
#[serde(default)]
pub diagnostics: Vec<String>,
#[serde(default)]
pub diagnostics_at: Option<SystemTime>,
/// 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<String>,
dropped: u64,
}
impl<'de> Deserialize<'de> for BoundedDiagnostics {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
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<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
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::<IgnoredAny>()?.is_none() {
break;
}
bounded.dropped = bounded.dropped.saturating_add(1);
continue;
}
let Some(diagnostic) = sequence.next_element::<Cow<'de, str>>()? 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<SystemTime>,
#[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<D>(deserializer: D) -> Result<Self, D::Error>
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 {
/// Record one observation where the logical and physical dimensions are
/// the same. Kept as a small compatibility helper for callers that only
/// have one size value.
pub fn record(&mut self, tier: &str, bytes: u64, versions: u64, objects: u64) {
self.record_dimensions(tier, bytes, bytes, versions, objects);
}
/// Record one observation without conflating logical usage with physical
/// 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 = 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_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| diagnostics_expired(at, now)) {
self.diagnostics.clear();
}
self.diagnostics_at = Some(now);
if self.diagnostics.iter().any(|entry| entry == &digest) {
return;
}
let current_bytes: usize = self.diagnostics.iter().map(String::len).sum();
if self.diagnostics.len() >= UNKNOWN_TIER_DIAGNOSTIC_ENTRY_CAP
|| current_bytes.saturating_add(digest.len()) > UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP
{
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 = 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.is_some_and(|at| diagnostics_expired(at, now)) {
self.diagnostics.clear();
self.diagnostics_at = None;
}
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());
}
}
if !self.diagnostics.is_empty() {
self.diagnostics_at = Some(now);
}
}
pub fn fits_add(&self, other: &Self) -> bool {
!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)
.is_some()
&& self.unknown_objects.checked_add(other.unknown_objects).is_some()
&& self.unknown_versions.checked_add(other.unknown_versions).is_some()
&& self.diagnostics_dropped.checked_add(other.diagnostics_dropped).is_some()
}
pub fn checked_add(&self, other: &Self) -> Option<Self> {
if !self.fits_add(other) {
return None;
}
let mut merged = self.clone();
merged.merge(other);
Some(merged)
}
pub fn is_empty(&self) -> bool {
!self.counter_overflowed
&& self.unknown_bytes == 0
&& self.unknown_physical_bytes == 0
&& self.unknown_objects == 0
&& self.unknown_versions == 0
&& self.diagnostics_dropped == 0
&& self.diagnostics.is_empty()
}
}
#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct AllTierStats {
pub tiers: HashMap<String, TierStats>,
@@ -124,6 +423,31 @@ impl AllTierStats {
.iter()
.all(|(tier, right)| self.tiers.get(tier).is_none_or(|left| left.fits_add(right)))
}
/// Fold keys from an older cache that are no longer present in the
/// current registry into the fixed unknown bucket. Built-in storage
/// classes remain known even when no remote tier is configured.
pub fn fold_unknown_tiers<'a, I>(&mut self, known_tiers: I)
where
I: IntoIterator<Item = &'a str>,
{
let known: HashSet<&str> = known_tiers.into_iter().collect();
let mut unknown = self.tiers.remove(UNKNOWN_TIER).unwrap_or_default();
let retired_tiers: Vec<String> = 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_tiers {
if let Some(stats) = self.tiers.remove(&tier) {
unknown = unknown.add(&stats);
}
}
if !unknown.is_empty() {
self.tiers.insert(UNKNOWN_TIER.to_string(), unknown);
}
}
}
/// Bucket target usage info provides replication statistics
@@ -211,6 +535,10 @@ pub struct DataUsageInfo {
/// tier exists, so an absent value means "not accounted", never "zero".
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tier_stats: Option<AllTierStats>,
/// Bounded diagnostics and separate logical/physical counters for objects
/// classified into [`UNKNOWN_TIER`].
#[serde(default, skip_serializing_if = "Option::is_none")]
pub unknown_tier_stats: Option<UnknownTierStats>,
/// Total number of buckets in this cluster
pub buckets_count: u64,
@@ -307,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<Self> {
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 {
@@ -336,6 +706,10 @@ pub struct SizeSummary {
pub repl_target_stats: HashMap<String, ReplTargetSizeSummary>,
/// Per-tier accounting, keyed by storage class or remote tier name
pub tier_stats: HashMap<String, TierStats>,
/// 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
@@ -427,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
@@ -541,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
@@ -681,6 +1063,13 @@ pub struct DataUsageEntry {
/// observed tier-classified objects.
#[serde(default)]
pub all_tier_stats: Option<AllTierStats>,
/// Bounded unknown-tier reconciliation state for this cache entry.
#[serde(default)]
pub unknown_tier_stats: Option<UnknownTierStats>,
/// Optional conservation proof. Missing values are legacy/unproven, not
/// zero-valued evidence.
#[serde(default)]
pub tier_accounting_proof: Option<TierAccountingProof>,
}
impl Serialize for DataUsageEntry {
@@ -691,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(11))?;
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)?;
@@ -703,11 +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)?;
// 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;
@@ -716,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;
@@ -744,6 +1161,29 @@ impl DataUsageEntry {
if let Some(o_tiers) = other.all_tier_stats.as_ref().filter(|tiers| !tiers.is_empty()) {
self.all_tier_stats.get_or_insert_with(AllTierStats::new).merge(o_tiers);
}
if let Some(other_unknown) = other.unknown_tier_stats.as_ref() {
self.unknown_tier_stats
.get_or_insert_with(UnknownTierStats::default)
.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);
@@ -757,6 +1197,15 @@ impl DataUsageEntry {
self.all_tier_stats.get_or_insert_with(AllTierStats::new).add_sizes(tiers);
}
pub fn add_unknown_tier_stats(&mut self, stats: &UnknownTierStats) {
if stats.is_empty() {
return;
}
self.unknown_tier_stats
.get_or_insert_with(UnknownTierStats::default)
.merge(stats);
}
pub fn checked_merge(&mut self, other: &DataUsageEntry) -> bool {
let scalar_counts_fit = self.objects.checked_add(other.objects).is_some()
&& self.versions.checked_add(other.versions).is_some()
@@ -820,8 +1269,21 @@ impl DataUsageEntry {
(_, None) | (None, Some(_)) => true,
(Some(left), Some(right)) => left.fits_merge(right),
};
let unknown_tier_stats_fit = match (&self.unknown_tier_stats, &other.unknown_tier_stats) {
(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 {
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);
@@ -1339,6 +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.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,
@@ -1766,6 +2229,20 @@ impl SizeSummary {
self.replica_count = self.replica_count.saturating_add(other.replica_count);
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;
// dropping this map here would recreate the original silent-loss bug
// at the cross-disk merge boundary.
for (tier, stats) in &other.tier_stats {
let entry = self.tier_stats.entry(tier.clone()).or_default();
*entry = entry.add(stats);
}
// Merge replication target stats
for (target, stats) in &other.repl_target_stats {
@@ -1878,6 +2355,59 @@ mod tests {
);
}
#[test]
fn retired_tier_stats_fold_into_fixed_unknown_bucket() {
let mut stats = AllTierStats::default();
stats.tiers.insert(
"RETIRED".to_string(),
TierStats {
total_size: 9,
num_versions: 2,
num_objects: 1,
},
);
stats.tiers.insert(
"WARM".to_string(),
TierStats {
total_size: 4,
num_versions: 1,
num_objects: 1,
},
);
stats.fold_unknown_tiers(["WARM"]);
assert!(!stats.tiers.contains_key("RETIRED"));
assert_eq!(stats.tiers.get("WARM").map(|v| v.total_size), Some(4));
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<String>,
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(
@@ -1901,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)]
@@ -1940,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 {
@@ -2708,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 {
+7 -5
View File
@@ -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,
};
}
+68 -13
View File
@@ -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<bool> = 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<Vec<u8>> {
fn generate_ns_scanner_capability_proof(
secret: &str,
challenge: Uuid,
server_epoch: Uuid,
supports_tier_registry_generation: bool,
) -> std::io::Result<Vec<u8>> {
if challenge.is_nil() || server_epoch.is_nil() {
return Err(std::io::Error::other("Invalid namespace scanner capability scope"));
}
let mut mac =
<HmacSha256 as KeyInit>::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 =
<HmacSha256 as KeyInit>::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<Vec<u8>> {
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<Vec<u8>> {
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,
@@ -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<Uuid> {
Err(Error::MethodNotAllowed)
}
async fn probe_ns_scanner_capability(&self, request: NsScannerCapabilityRequest) -> Result<NsScannerCapabilityResponse> {
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<Uuid> {
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<NsScannerCapabilityResponse> {
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<NsScannerCapabilityResponse> {
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<Option<Uuid>> {
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();
+3 -2
View File
@@ -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,
+58 -3
View File
@@ -781,14 +781,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)
@@ -4040,10 +4042,21 @@ mod tests {
NsScannerProbe(NsScannerCapabilityRequest),
}
#[derive(Debug, Clone, Default)]
#[derive(Debug, Clone)]
struct RecordingInternodeDataTransport {
calls: Arc<StdMutex<Vec<RecordedTransportCall>>>,
ns_scanner_probe_status: Arc<StdMutex<Option<u16>>>,
ns_scanner_generation_support: Arc<StdMutex<Option<bool>>>,
}
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)]
@@ -4263,6 +4276,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<bool>) -> Self {
Self {
calls: Arc::default(),
ns_scanner_probe_status: Arc::default(),
ns_scanner_generation_support: Arc::new(StdMutex::new(support)),
}
}
@@ -4945,6 +4967,23 @@ mod tests {
Ok(Uuid::from_u128(1))
}
async fn probe_ns_scanner_capability(
&self,
request: NsScannerCapabilityRequest,
) -> Result<crate::storage_api_contracts::internode::NsScannerCapabilityResponse> {
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"
}
@@ -6757,6 +6796,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;
@@ -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,
};
}
+133 -12
View File
@@ -29,7 +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, 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};
@@ -205,25 +206,79 @@ impl ScannerSizeSummaryExt for SizeSummary {
self.versions = self.versions.saturating_add(1);
}
let size = usize::try_from(size.max(0)).unwrap_or(usize::MAX);
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));
// 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;
}
if let Some(tier_stats) = self.tier_stats.get_mut(&tier) {
*tier_stats = tier_stats.add(&TierStats {
total_size: u64::try_from(oi.size).unwrap_or(0),
num_versions: 1,
num_objects: u64::from(oi.is_latest),
});
// 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_bytes,
num_versions: 1,
num_objects: 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);
}
}
@@ -271,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<u64>,
}
#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq, Hash)]
@@ -344,6 +404,10 @@ pub struct DataUsageCacheInfo {
pub scan_plan_digest: Option<DataUsageScanPlanDigest>,
#[serde(default)]
pub cache_key_format: u16,
/// Registry generation used for the completed/partial scan. This is
/// process-local audit data; older cache writers omit it.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tier_registry_generation: Option<u64>,
}
impl Serialize for DataUsageCacheInfo {
@@ -353,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(16))?;
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)?;
@@ -370,6 +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)?;
if let Some(generation) = self.tier_registry_generation {
state.serialize_entry("tier_registry_generation", &generation)?;
}
state.end()
}
}
@@ -390,6 +458,58 @@ pub(crate) enum DataUsageCachePrepareOutcome {
}
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::<HashSet<_>>();
for entry in self.cache.values_mut() {
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;
}
}
}
/// Prefix-level usage query over this (writer-side) cache; see
/// [`prefix_usage_in_cache`] for the semantics
/// (rustfs/backlog#1872).
@@ -875,6 +995,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.filter(|stats| !stats.is_empty()),
buckets_count: u64::try_from(buckets.len()).unwrap_or(u64::MAX),
buckets_usage,
..Default::default()
+263 -1
View File
@@ -573,7 +573,6 @@ fn size_summary_add_saturates_all_usage_counters() {
failed_count: usize::MAX,
},
);
let mut increment = SizeSummary {
total_size: 1,
versions: 1,
@@ -588,6 +587,24 @@ fn size_summary_add_saturates_all_usage_counters() {
failed_count: 1,
..Default::default()
};
summary.tier_stats.insert(
UNKNOWN_TIER.to_string(),
TierStats {
total_size: u64::MAX,
num_versions: u64::MAX,
num_objects: u64::MAX,
},
);
increment.tier_stats.insert(
UNKNOWN_TIER.to_string(),
TierStats {
total_size: 1,
num_versions: 1,
num_objects: 1,
},
);
increment.unknown_tier_stats.unknown_bytes = 1;
increment.unknown_tier_stats.unknown_physical_bytes = 1;
increment.repl_target_stats.insert(
target.clone(),
ReplTargetSizeSummary {
@@ -624,6 +641,8 @@ fn size_summary_add_saturates_all_usage_counters() {
assert_eq!(target_summary.failed_size, i64::MAX);
assert_eq!(target_summary.pending_count, usize::MAX);
assert_eq!(target_summary.failed_count, usize::MAX);
assert_eq!(summary.tier_stats[UNKNOWN_TIER].total_size, u64::MAX);
assert_eq!(summary.unknown_tier_stats.unknown_bytes, 1);
}
#[test]
@@ -673,6 +692,249 @@ fn size_summary_actions_accounting_accumulates_tier_stats() {
);
}
#[test]
fn unknown_tier_is_bounded_and_accounted() {
let mut summary = SizeSummary::new();
summary.tier_stats.insert("WARM".to_string(), TierStats::default());
let object = ObjectInfo {
storage_class: Some("retired-tier".to_string()),
size: 11,
is_latest: true,
..Default::default()
};
summary.actions_accounting(&object, 11, 11);
assert_eq!(summary.tier_stats.len(), 2);
assert_eq!(summary.tier_stats.get(UNKNOWN_TIER).map(|stats| stats.total_size), Some(11));
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::<usize>() <= UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP);
assert!(
summary
.unknown_tier_stats
.diagnostics
.iter()
.all(|entry| !entry.contains("retired"))
);
}
#[test]
fn unknown_tier_is_accounted_when_no_remote_tier_is_configured() {
let mut summary = SizeSummary::new();
let object = ObjectInfo {
storage_class: Some("retired-tier".to_string()),
size: 3,
is_latest: true,
..Default::default()
};
summary.actions_accounting(&object, 9, 9);
assert_eq!(summary.tier_stats.len(), 1);
assert_eq!(summary.tier_stats[UNKNOWN_TIER].total_size, 3);
assert_eq!(summary.unknown_tier_stats.unknown_bytes, 9);
assert_eq!(summary.unknown_tier_stats.unknown_physical_bytes, 3);
let standard = ObjectInfo {
storage_class: Some(storageclass::STANDARD.to_string()),
size: 4,
is_latest: true,
..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");
}
#[test]
fn million_unique_tier_keys_do_not_grow_stats_map() {
let mut summary = SizeSummary::new();
summary.tier_stats.insert("WARM".to_string(), TierStats::default());
for index in 0..1_000_000_u64 {
let object = ObjectInfo {
storage_class: Some(format!("untrusted-tier-{index}")),
size: 1,
..Default::default()
};
summary.actions_accounting(&object, 1, 1);
}
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::<usize>() <= UNKNOWN_TIER_DIAGNOSTIC_BYTE_CAP);
}
#[test]
fn unknown_tier_never_triggers_transition() {
let mut summary = SizeSummary::new();
summary.tier_stats.insert("WARM".to_string(), TierStats::default());
let mut object = ObjectInfo {
storage_class: Some("removed-tier".to_string()),
size: 7,
..Default::default()
};
object.transitioned_object.status = TRANSITION_COMPLETE.to_string();
object.transitioned_object.tier = "removed-tier".to_string();
summary.actions_accounting(&object, 7, 7);
assert_eq!(summary.tier_stats.get("removed-tier"), None);
assert_eq!(summary.tier_stats[UNKNOWN_TIER].total_size, 7);
}
#[test]
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("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(&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
);
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]
fn tier_registry_refresh_does_not_mix_cycle_generations() {
let first = crate::TierRegistrySnapshot {
generation: 1,
names: Arc::from(["WARM".to_string()]),
refresh_failed: false,
};
let second = crate::TierRegistrySnapshot {
generation: 2,
names: Arc::from(["COLD".to_string()]),
refresh_failed: false,
};
assert_ne!(first.generation, second.generation);
assert_eq!(first.names.as_ref(), ["WARM".to_string()]);
assert_eq!(second.names.as_ref(), ["COLD".to_string()]);
assert!(first.refreshed(Err(())).refresh_failed);
assert!(!second.refreshed(Ok(Arc::from(["HOT".to_string()]))).refresh_failed);
assert_eq!(first.refreshed(Err(())).generation, first.generation);
assert_eq!(first.refreshed(Err(())).names, first.names);
}
#[test]
fn unknown_tier_counter_uses_checked_arithmetic() {
let max = TierStats {
total_size: u64::MAX,
num_versions: u64::MAX,
num_objects: u64::MAX,
};
assert!(max.checked_add(&TierStats::default()).is_some());
assert!(
max.checked_add(&TierStats {
total_size: 1,
..Default::default()
})
.is_none()
);
let mut unknown = UnknownTierStats {
unknown_bytes: u64::MAX,
..Default::default()
};
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]
fn test_data_usage_entry_merge_sums_failed_objects() {
let mut left = DataUsageEntry {
+286 -10
View File
@@ -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::{
@@ -94,6 +97,45 @@ static SCANNER_RUNTIME_INSTANCES: AtomicU64 = AtomicU64::new(0);
static SCANNER_FOREGROUND_READ_ACTIVITY: AtomicU64 = AtomicU64::new(0);
static SCANNER_FOREGROUND_STREAM_READS: AtomicU64 = AtomicU64::new(0);
/// Immutable tier registry captured at the beginning of a folder scan.
/// Generation makes it possible to prove that a result was classified against
/// one registry even when the process-wide TTL cache refreshes concurrently.
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct TierRegistrySnapshot {
pub(crate) generation: u64,
pub(crate) names: Arc<[String]>,
/// True when the last refresh attempt failed and `names` is therefore a
/// retained last-good snapshot rather than a newly read registry.
pub(crate) refresh_failed: bool,
}
impl TierRegistrySnapshot {
/// Apply a refresh only when the registry read succeeds. A failed refresh
/// retains the prior generation, preventing a transient config failure
/// from classifying the remainder of a scan against an empty registry.
pub(crate) fn refreshed(&self, names: Result<Arc<[String]>, ()>) -> Self {
match names {
Ok(names) => Self {
generation: self.generation.saturating_add(1),
names,
refresh_failed: false,
},
Err(()) => Self {
refresh_failed: true,
..self.clone()
},
}
}
pub(crate) fn initial(names: Arc<[String]>) -> Self {
Self {
generation: 1,
names,
refresh_failed: false,
}
}
}
pub fn current_scanner_activity() -> u64 {
SCANNER_ACTIVE_WORK_UNITS.load(Ordering::Relaxed)
}
@@ -371,6 +413,7 @@ pub(crate) fn resolve_scanner_server_config() -> Option<ServerConfig> {
/// 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.
///
@@ -383,24 +426,192 @@ const TIER_NAME_CACHE_TTL: Duration = Duration::from_secs(30);
/// `TIER_NAME_CACHE_TTL` later; a removed tier can leave an all-zero
/// `TierStats` seed behind for one cache generation, which merges harmlessly
/// by key in per-object accounting and disappears on the next refresh.
static TIER_NAME_CACHE: RwLock<Option<(Instant, Arc<[String]>)>> = RwLock::new(None);
static TIER_NAME_CACHE: RwLock<Option<(Instant, TierRegistrySnapshot)>> = RwLock::new(None);
static TIER_REGISTRY_GENERATION: AtomicU64 = AtomicU64::new(0);
static TIER_CYCLE_SNAPSHOTS: LazyLock<RwLock<HashMap<(u64, u64), TierRegistrySnapshot>>> =
LazyLock::new(|| RwLock::new(HashMap::new()));
static TIER_ACTIVE_CYCLES: LazyLock<RwLock<HashMap<(u64, u64), usize>>> = LazyLock::new(|| RwLock::new(HashMap::new()));
static TIER_NAME_REFRESH_LOCK: LazyLock<tokio::sync::Mutex<()>> = 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 {
{
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;
}
}
// 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 names = ecstore_get_global_tier_config_mgr()
.read()
.await
.list_tiers()
.into_iter()
.map(|tier| tier.name)
.collect::<Vec<_>>();
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)
},
}
}
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<String>) -> Result<Arc<[String]>, ()> {
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_NAME_CACHE.read().unwrap_or_else(|err| err.into_inner()).clone();
if let Some((refreshed_at, names)) = cached
&& refreshed_at.elapsed() < TIER_NAME_CACHE_TTL
{
return names;
let cached = TIER_CYCLE_SNAPSHOTS.read().unwrap_or_else(|err| err.into_inner());
if let Some(snapshot) = cached.get(&key) {
return snapshot.clone();
}
}
let tiers = ecstore_get_global_tier_config_mgr().read().await.list_tiers();
let names: Arc<[String]> = tiers.iter().map(|tier| tier.name.clone()).collect::<Vec<_>>().into();
*TIER_NAME_CACHE.write().unwrap_or_else(|err| err.into_inner()) = Some((Instant::now(), Arc::clone(&names)));
names
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
@@ -408,6 +619,9 @@ pub(crate) async fn runtime_tier_names() -> Arc<[String]> {
#[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) {
@@ -616,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();
+36 -3
View File
@@ -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",
@@ -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();
+32 -4
View File
@@ -55,9 +55,9 @@ use tracing::{debug, error, warn};
use crate::{
Disk, DiskError, DiskInfoOptions, Evaluator, Event, LcEventSrc, ListPathRawOptions, ObjectOpts, ReplicationConfig,
ReplicationHealObject, ReplicationQueueAdmission, ReplicationStatusType, STORAGE_FORMAT_FILE, ScannerDiskExt as _,
ScannerLifecycleConfigExt as _, ScannerVersioningConfigExt as _, StorageError, 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, scanner_is_erasure,
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_for_cycle, scanner_is_erasure,
scanner_replication_config_for_lifecycle_eval,
};
use crate::{ScannerObjectInfo as ObjectInfo, ScannerObjectToDelete as ObjectToDelete};
@@ -626,6 +626,20 @@ 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 {
@@ -636,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 {
@@ -670,6 +686,9 @@ pub struct FolderScanner {
budget: Arc<ScannerCycleBudget>,
skip_heal: Arc<std::sync::atomic::AtomicBool>,
local_disk: Arc<Disk>,
/// Tier registry frozen for this folder scan. A refresh applies to the
/// next scan and cannot mix generations in one aggregate.
tier_registry: TierRegistrySnapshot,
pending_heals_changed: bool,
#[cfg(test)]
list_path_raw_options_observer: Option<mpsc::UnboundedSender<ListPathRawTimeoutSnapshot>>,
@@ -1329,7 +1348,11 @@ impl FolderScanner {
continue;
}
let sz = match self.local_disk.get_size(item.clone()).await {
let sz = match self
.local_disk
.get_size_with_tier_names(item.clone(), &self.tier_registry.names)
.await
{
Ok(sz) => sz,
Err(e) => {
let failure_action = classify_get_size_failure(&item, &e);
@@ -2161,6 +2184,10 @@ 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_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);
// Create folder scanner
let mut scanner = FolderScanner {
@@ -2189,6 +2216,7 @@ pub async fn scan_data_folder(
budget: budget.clone(),
skip_heal,
local_disk,
tier_registry,
pending_heals_changed: false,
#[cfg(test)]
list_path_raw_options_observer: None,
@@ -289,11 +289,62 @@ impl ScannerItem {
item.object_path()
}
fn effective_tier(oi: &ObjectInfo) -> &str {
if oi.transitioned_object.status == crate::TRANSITION_COMPLETE {
oi.transitioned_object.tier.as_str()
} else {
oi.storage_class.as_deref().unwrap_or(crate::storageclass::STANDARD)
}
}
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 {
matches!(
action,
IlmAction::TransitionAction
| IlmAction::TransitionVersionAction
| IlmAction::DeleteAction
| IlmAction::DeleteVersionAction
| IlmAction::DeleteRestoredAction
| IlmAction::DeleteRestoredVersionAction
| IlmAction::DeleteAllVersionsAction
| IlmAction::DelMarkerDeleteAllVersionsAction
)
}
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<ObjectInfo>,
lock_retention: Option<Arc<ObjectLockConfiguration>>,
versioning_config: VersioningConfiguration,
tier_names: &[String],
size_summary: &mut SizeSummary,
) {
let object_path = self.object_path();
@@ -402,6 +453,9 @@ impl ScannerItem {
let mut noncurrent_accounting: Vec<PendingScannerAccounting<'_>> = 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];
@@ -425,6 +479,18 @@ impl ScannerItem {
let mut size = actual_size;
let mut account_now = true;
// A retired/unknown source tier may point at a remote object
// that cannot be safely deleted or transitioned. Lifecycle
// evaluation is still useful for accounting, but all
// side-effecting tier actions fail closed until the registry
// recognizes the source again.
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;
continue;
}
match event.action {
IlmAction::DeleteAllVersionsAction | IlmAction::DelMarkerDeleteAllVersionsAction => {
debug!(
@@ -929,4 +995,49 @@ mod tests {
assert_eq!(item.object_name, "object");
assert_eq!(item.object_path(), "object");
}
#[test]
fn unknown_tier_never_triggers_transition() {
let object = ObjectInfo {
storage_class: Some("retired-tier".to_string()),
..Default::default()
};
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
));
}
}
@@ -325,6 +325,11 @@ async fn build_test_scanner() -> (FolderScanner, std::path::PathBuf) {
budget: ScannerCycleBudget::new(&CancellationToken::new(), Default::default()),
skip_heal: Arc::new(AtomicBool::new(false)),
local_disk: disk,
tier_registry: crate::TierRegistrySnapshot {
generation: 0,
names: Arc::new([]),
refresh_failed: false,
},
pending_heals_changed: false,
list_path_raw_options_observer: None,
};
+20 -4
View File
@@ -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<Vec<BucketInfo>>,
digest: DataUsageScanPlanDigest,
leader_epoch: u64,
tier_registry_generation: u64,
dirty_usage_buckets: Arc<DirtyUsageBuckets>,
bucket_failures: ScannerBucketFailureState,
pending_maintenance_work: Arc<AtomicBool>,
@@ -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 {
@@ -498,6 +508,9 @@ pub trait ScannerIODisk: Send + Sync + Debug + 'static {
) -> Result<ScannerDiskScanOutcome>;
async fn get_size(&self, item: ScannerItem) -> Result<SizeSummary>;
/// Read one object using a registry snapshot captured at scan start.
async fn get_size_with_tier_names(&self, item: ScannerItem, tier_names: &[String]) -> Result<SizeSummary>;
}
#[derive(Debug)]
@@ -649,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,
+181 -6
View File
@@ -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<u64>,
) -> std::result::Result<Option<DataUsageEntryInfo>, 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<u64>,
}
#[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<DataUsageCacheSource>,
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()
@@ -216,6 +286,7 @@ pub(super) fn completed_data_usage_info(
delete_markers_total_count: u64::try_from(total.delete_markers).ok()?,
objects_total_size: u64::try_from(total.size).ok()?,
tier_stats: total.all_tier_stats.filter(|tiers| !tiers.is_empty()),
unknown_tier_stats: total.unknown_tier_stats.filter(|stats| !stats.is_empty()),
buckets_count: u64::try_from(all_buckets.len()).ok()?,
bucket_sizes,
buckets_usage,
@@ -225,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<u64>,
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<DataUsageEntryInfo>,
cache: &DataUsageCache,
@@ -319,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(_))
) {
+20 -4
View File
@@ -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),
+11
View File
@@ -48,6 +48,7 @@ impl ScannerIOCycle for ECStore {
scan_mode: HealScanMode,
) -> Result<ScannerCycleResult> {
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 {
+27 -15
View File
@@ -13,29 +13,38 @@
// limitations under the License.
/// ScannerIODisk implementation for Disk: get_size and the per-disk bucket scan.
use super::*;
use crate::UNKNOWN_TIER;
///
/// Seed [`SizeSummary::tier_stats`] from the cached tier-name list.
///
/// Preserves the original seeding semantics: with no tiers configured the map
/// stays completely empty (STANDARD/RRS are not seeded either); otherwise the
/// standard storage classes are seeded alongside every configured tier so
/// per-object accounting always finds its tier key.
/// Preserves the original no-tier shape: with no tiers configured the map
/// stays completely empty (STANDARD/RRS/UNKNOWN are not seeded either).
/// Otherwise the standard storage classes and one fixed unknown bucket are
/// seeded alongside every configured tier so per-object accounting never
/// inserts an untrusted metadata key.
pub(super) fn tier_stats_template(tier_names: &[String]) -> HashMap<String, TierStats> {
let mut tier_stats = HashMap::with_capacity(tier_names.len() + 2);
let mut tier_stats = HashMap::with_capacity(tier_names.len() + 3);
for tier_name in tier_names {
tier_stats.insert(tier_name.clone(), TierStats::default());
if tier_name != UNKNOWN_TIER {
tier_stats.insert(tier_name.clone(), TierStats::default());
}
}
if !tier_stats.is_empty() {
tier_stats.insert(storageclass::STANDARD.to_string(), TierStats::default());
tier_stats.insert(storageclass::RRS.to_string(), TierStats::default());
tier_stats.insert(UNKNOWN_TIER.to_string(), TierStats::default());
}
tier_stats
}
#[async_trait::async_trait]
impl ScannerIODisk for Disk {
async fn get_size(&self, mut item: ScannerItem) -> Result<SizeSummary> {
async fn get_size(&self, item: ScannerItem) -> Result<SizeSummary> {
self.get_size_with_tier_names(item, &runtime_tier_names().await).await
}
async fn get_size_with_tier_names(&self, mut item: ScannerItem, tier_names: &[String]) -> Result<SizeSummary> {
let done_object = Metrics::time(Metric::ScanObject);
if !is_xl_meta_path(&item.path) {
@@ -105,12 +114,13 @@ impl ScannerIODisk for Disk {
.map(|v| ObjectInfo::from_file_info(v, item.bucket.as_str(), object_path.as_str(), versioned))
.collect::<Vec<ObjectInfo>>();
let mut size_summary = SizeSummary::default();
// Tier names come from the process-wide TTL cache; seeding from them
// replaces the per-object clone of every full TierConfig.
let tier_names = runtime_tier_names().await;
size_summary.tier_stats = tier_stats_template(&tier_names);
// 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.
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;
@@ -120,12 +130,14 @@ impl ScannerIODisk for Disk {
// `object_infos`.
global_metrics().record_scanner_versions_scanned(object_infos.len() as u64);
item.apply_actions(object_infos, lock_config, versioning_config, &mut size_summary)
item.apply_actions(object_infos, lock_config, versioning_config, tier_names, &mut size_summary)
.await;
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;
}
}
}
@@ -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::<HashSet<_>>();
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]
+44 -5
View File
@@ -23,7 +23,7 @@ use crate::storage_api::owner::{
use crate::storage_api::scan::{BucketOperations as _, DeleteBucketOptions, MakeBucketOptions, ObjectIO as _};
use crate::{
DiskOption, ECStore, Endpoint, EndpointServerPools, Endpoints, InstanceContext, PoolEndpoints, ScannerObjectOptions,
ScannerPutObjReader, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
ScannerPutObjReader, UNKNOWN_TIER, init_bucket_metadata_sys_for_scanner_tests, init_ecstore_config_for_scanner_tests,
init_local_disks_with_instance_ctx, new_disk, path2_bucket_object_with_base_path,
};
use rustfs_filemeta::FileInfo;
@@ -986,8 +986,8 @@ fn is_xl_meta_path_accepts_forward_separator() {
fn tier_stats_template_seeds_tiers_and_standard_classes() {
let template = tier_stats_template(&["WARM".to_string(), "COLD".to_string()]);
assert_eq!(template.len(), 4);
for tier in ["WARM", "COLD", storageclass::STANDARD, storageclass::RRS] {
assert_eq!(template.len(), 5);
for tier in ["WARM", "COLD", storageclass::STANDARD, storageclass::RRS, UNKNOWN_TIER] {
assert_eq!(template.get(tier), Some(&TierStats::default()), "missing seed for tier {tier}");
}
}
@@ -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());
}
+3
View File
@@ -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<u8>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub supports_tier_registry_generation: Option<bool>,
}
pub mod admin;
+4 -11
View File
@@ -158,11 +158,10 @@ added by backlog#1153 infra-4.
## Coverage
Workspace line coverage is measured weekly. Pull requests that touch iam, kms,
policy, or crypto also run a non-required, report-only comparison against
`.config/coverage-baselines.toml`. During calibration, a regression is recorded
in the job summary without failing the job; missing or malformed coverage
evidence still fails closed (backlog#1153 infra-6).
Line coverage is measured **weekly, not per-PR**, and is non-blocking: it
exists for visibility and trend, never as a required check. Per-crate ratchets
for the security-critical crates (iam / kms / policy / crypto) build on this
baseline later (backlog#1153 infra-6, report-only first).
- **CI**: `.github/workflows/coverage.yml` runs every Sunday and on manual
dispatch: `cargo llvm-cov nextest --workspace --exclude e2e_test` under the
@@ -175,12 +174,6 @@ evidence still fails closed (backlog#1153 infra-6).
plus the full suite). It prints the same per-crate table via
`scripts/coverage_per_crate.py` and writes `target/llvm-cov/lcov.info` and
`coverage.json`.
- **Security-critical ratchet**: relevant pull requests compare iam / kms /
policy / crypto line coverage with the versioned baseline. Drops greater than
the configured one-percentage-point calibration threshold are marked
`REGRESSION (report-only)`. The weekly summary runs the same comparison so
calibration continues even when no relevant pull request is open. Baseline
changes require a linked coverage run and a reviewed explanation.
- **Trend comparison**: each run's job summary is the weekly per-crate
snapshot — open two runs from the Actions history (workflow "coverage") and
compare their tables. For line-level diffs, download the two runs'
+84 -29
View File
@@ -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<u16>,
ns_scanner_challenge: Option<uuid::Uuid>,
#[serde(rename = "ns_scanner_tier_registry_generation")]
ns_scanner_tier_registry_generation: Option<bool>,
}
fn verify_ns_scanner_body_digest(query: &NsScannerQuery, body: &[u8]) -> bool {
@@ -410,7 +414,9 @@ async fn handle_internode_rpc(req: Request<Incoming>) -> Response<Body> {
(Method::GET, WALK_DIR_PATH) | (Method::HEAD, WALK_DIR_PATH) => handle_walk_dir(req).await,
(Method::GET, NS_SCANNER_PATH) => match parse_query::<NsScannerCapabilityQuery>(&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<Body> {
fn ns_scanner_capability_response(challenge: uuid::Uuid, include_tier_registry_generation: bool) -> Response<Body> {
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::<NsScannerQuery>(&query).is_err());
assert!(serde_urlencoded::from_str::<super::NsScannerCapabilityQuery>("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::<super::NsScannerCapabilityQuery>(&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<u8>,
}
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(&current).expect("new response should decode");
assert_eq!(decoded.supports_tier_registry_generation, Some(true));
assert!(rmp_serde::from_slice::<LegacyCapabilityResponse>(&current).is_err());
}
#[test]
+18 -9
View File
@@ -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)]
@@ -1731,8 +1732,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<Vec<u8>> {
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<Vec<u8>> {
ecstore_rpc::sign_ns_scanner_capability_with_tier_registry_generation(
challenge,
server_epoch,
supports_tier_registry_generation,
)
}
pub(crate) fn sign_put_file_capability(
-264
View File
@@ -1,264 +0,0 @@
#!/usr/bin/env python3
# Copyright 2024 RustFS Team
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Compare security-critical crate line coverage with the report-only baseline."""
import argparse
import json
import math
import os
import sys
import tempfile
import tomllib
from pathlib import Path
from coverage_per_crate import fmt_pct, load_coverage
SECURITY_CRATES = ("crates/iam", "crates/kms", "crates/policy", "crates/crypto")
def load_baselines(path: str) -> tuple[float, dict[str, tuple[int, int]]]:
with open(path, "rb") as fh:
config = tomllib.load(fh)
if config.get("phase") != "report-only":
raise ValueError("coverage baseline phase must be report-only")
raw_allowed_drop = config["allowed_drop_percentage_points"]
if isinstance(raw_allowed_drop, bool) or not isinstance(raw_allowed_drop, (int, float)):
raise ValueError("allowed_drop_percentage_points must be a number")
allowed_drop = float(raw_allowed_drop)
if not math.isfinite(allowed_drop) or allowed_drop < 0:
raise ValueError("allowed_drop_percentage_points must be finite and non-negative")
baselines: dict[str, tuple[int, int]] = {}
for crate, values in config["crates"].items():
covered = values["covered"]
count = values["count"]
if type(covered) is not int or type(count) is not int:
raise ValueError(f"invalid baseline for {crate}: covered and count must be integers")
if covered < 0 or count <= 0 or covered > count:
raise ValueError(f"invalid baseline for {crate}: {covered}/{count}")
baselines[crate] = (covered, count)
missing = [crate for crate in SECURITY_CRATES if crate not in baselines]
unexpected = sorted(set(baselines).difference(SECURITY_CRATES))
if missing or unexpected:
raise ValueError(f"coverage baseline crate set mismatch: missing={missing}, unexpected={unexpected}")
return allowed_drop, baselines
def compare(
current: dict[str, list[int]],
baselines: dict[str, tuple[int, int]],
allowed_drop: float,
) -> list[tuple[str, int, int, int, int, float, bool]]:
rows = []
for crate, (baseline_covered, baseline_count) in baselines.items():
if crate not in current:
raise ValueError(f"coverage report is missing {crate}")
covered, count = current[crate]
if type(covered) is not int or type(count) is not int:
raise ValueError(f"invalid coverage for {crate}: covered and count must be integers")
if covered < 0 or count <= 0 or covered > count:
raise ValueError(f"invalid coverage for {crate}: {covered}/{count}")
current_pct = 100.0 * covered / count
baseline_pct = 100.0 * baseline_covered / baseline_count
delta = current_pct - baseline_pct
rows.append((crate, covered, count, baseline_covered, baseline_count, delta, delta < -allowed_drop))
return rows
def print_report(rows: list[tuple[str, int, int, int, int, float, bool]], allowed_drop: float) -> None:
print("## Security-critical coverage ratchet (report-only)")
print()
print(f"Calibration threshold: a drop greater than {allowed_drop:.2f} percentage points is reported as a regression.")
print()
print("| Crate | Current | Baseline | Delta | Status |")
print("|---|---:|---:|---:|---|")
for crate, covered, count, baseline_covered, baseline_count, delta, regressed in rows:
status = "REGRESSION (report-only)" if regressed else "OK"
print(
f"| `{crate}` | {fmt_pct(covered, count)} ({covered}/{count}) "
f"| {fmt_pct(baseline_covered, baseline_count)} ({baseline_covered}/{baseline_count}) "
f"| {delta:+.2f} pp | {status} |"
)
print()
print("This calibration phase records regressions without failing the job; malformed or incomplete evidence still fails closed.")
def self_test() -> None:
with tempfile.TemporaryDirectory() as tmp:
root = Path(tmp)
coverage = root / "coverage.json"
baseline = root / "baseline.toml"
coverage_data = {
"data": [
{
"files": [
{
"filename": str(root / "crates/iam/src/lib.rs"),
"summary": {"lines": {"covered": 80, "count": 100}},
},
{
"filename": str(root / "crates/kms/src/lib.rs"),
"summary": {"lines": {"covered": 90, "count": 100}},
},
{
"filename": str(root / "crates/policy/src/lib.rs"),
"summary": {"lines": {"covered": 90, "count": 100}},
},
{
"filename": str(root / "crates/crypto/src/lib.rs"),
"summary": {"lines": {"covered": 90, "count": 100}},
},
],
"totals": {"lines": {"covered": 350, "count": 400}},
}
]
}
coverage.write_text(json.dumps(coverage_data), encoding="utf-8")
baseline_text = """phase = "report-only"
allowed_drop_percentage_points = 1.0
[crates."crates/iam"]
covered = 90
count = 100
[crates."crates/kms"]
covered = 85
count = 100
[crates."crates/policy"]
covered = 90
count = 100
[crates."crates/crypto"]
covered = 90
count = 100
"""
baseline.write_text(baseline_text, encoding="utf-8")
current, _ = load_coverage(str(coverage), str(root))
allowed_drop, baselines = load_baselines(str(baseline))
rows = compare(current, baselines, allowed_drop)
assert [row[-1] for row in rows] == [True, False, False, False]
try:
compare({"crates/iam": current["crates/iam"]}, baselines, allowed_drop)
except ValueError as error:
assert str(error) == "coverage report is missing crates/kms"
else:
raise AssertionError("missing crate must fail closed")
try:
compare({**current, "crates/iam": [101, 100]}, baselines, allowed_drop)
except ValueError as error:
assert str(error) == "invalid coverage for crates/iam: 101/100"
else:
raise AssertionError("invalid coverage must fail closed")
for invalid_threshold in ("true", '"1.0"', "nan", "inf", "-inf"):
baseline.write_text(
baseline_text.replace("allowed_drop_percentage_points = 1.0", f"allowed_drop_percentage_points = {invalid_threshold}"),
encoding="utf-8",
)
try:
load_baselines(str(baseline))
except ValueError:
pass
else:
raise AssertionError(f"non-finite threshold {invalid_threshold} must fail closed")
for field, invalid_values in (
("covered", ("true", '"90"', "90.0", "90.5")),
("count", ("true", '"100"', "100.0", "100.5")),
):
for invalid_value in invalid_values:
baseline.write_text(
baseline_text.replace(f"{field} = {90 if field == 'covered' else 100}", f"{field} = {invalid_value}", 1),
encoding="utf-8",
)
try:
load_baselines(str(baseline))
except ValueError:
pass
else:
raise AssertionError(f"non-integer baseline {field} {invalid_value} must fail closed")
for covered, count in (
(True, 100),
(80, True),
(80.0, 100),
(80, 100.0),
(float("nan"), 100),
(80, float("inf")),
):
try:
compare({**current, "crates/iam": [covered, count]}, baselines, allowed_drop)
except ValueError:
pass
else:
raise AssertionError(f"invalid aggregate coverage {covered}/{count} must fail closed")
lines = coverage_data["data"][0]["files"][0]["summary"]["lines"]
for field, invalid_values in (
("covered", (True, "80", 80.0, 80.5, float("nan"), float("inf"), float("-inf"))),
("count", (True, "100", 100.0, 100.5, float("nan"), float("inf"), float("-inf"))),
):
original = lines[field]
for invalid_value in invalid_values:
lines[field] = invalid_value
coverage.write_text(json.dumps(coverage_data), encoding="utf-8")
try:
load_coverage(str(coverage), str(root))
except ValueError:
pass
else:
raise AssertionError(f"invalid raw coverage {field} {invalid_value} must fail closed")
lines[field] = original
baseline.write_text(
baseline_text.replace(
'[crates."crates/crypto"]\ncovered = 90\ncount = 100\n',
"",
),
encoding="utf-8",
)
try:
load_baselines(str(baseline))
except ValueError:
pass
else:
raise AssertionError("missing security-crate baseline must fail closed")
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("coverage_json", nargs="?")
parser.add_argument("--baseline", default=".config/coverage-baselines.toml")
parser.add_argument("--repo-root", default=os.getcwd())
parser.add_argument("--self-test", action="store_true")
args = parser.parse_args()
if args.self_test:
self_test()
print("security coverage self-test passed")
return 0
if not args.coverage_json:
parser.error("coverage_json is required unless --self-test is used")
try:
current, _ = load_coverage(args.coverage_json, os.path.abspath(args.repo_root))
allowed_drop, baselines = load_baselines(args.baseline)
rows = compare(current, baselines, allowed_drop)
except (OSError, ValueError, KeyError, IndexError, json.JSONDecodeError, tomllib.TOMLDecodeError) as error:
print(f"error: {error}", file=sys.stderr)
return 1
print_report(rows, allowed_drop)
return 0
if __name__ == "__main__":
sys.exit(main())
+14 -27
View File
@@ -47,31 +47,6 @@ def fmt_pct(covered: int, count: int) -> str:
return f"{100.0 * covered / count:.2f}%" if count else ""
def _line_counts(lines: dict[str, int], source: str) -> tuple[int, int]:
covered = lines["covered"]
count = lines["count"]
if type(covered) is not int or type(count) is not int or covered < 0 or count < 0 or covered > count:
raise ValueError(f"invalid line coverage for {source}: {covered}/{count}")
return covered, count
def load_coverage(path: str, root: str) -> tuple[dict[str, list[int]], dict[str, int]]:
with open(path, encoding="utf-8") as fh:
export = json.load(fh)
data = export["data"][0]
files = data["files"]
total_covered, total_count = _line_counts(data["totals"]["lines"], "totals")
crates: dict[str, list[int]] = {}
for f in files:
covered, count = _line_counts(f["summary"]["lines"], f["filename"])
acc = crates.setdefault(crate_label(f["filename"], root), [0, 0])
acc[0] += covered
acc[1] += count
return crates, {"covered": total_covered, "count": total_count}
def main() -> int:
if len(sys.argv) < 2 or len(sys.argv) > 3:
print(__doc__.strip(), file=sys.stderr)
@@ -79,12 +54,24 @@ def main() -> int:
path = sys.argv[1]
root = os.path.abspath(sys.argv[2] if len(sys.argv) == 3 else os.getcwd())
with open(path, encoding="utf-8") as fh:
export = json.load(fh)
try:
crates, totals = load_coverage(path, root)
except (KeyError, IndexError, ValueError) as exc:
data = export["data"][0]
files = data["files"]
totals = data["totals"]["lines"]
except (KeyError, IndexError) as exc:
print(f"error: unexpected llvm-cov JSON shape ({exc})", file=sys.stderr)
return 1
crates: dict[str, list[int]] = {}
for f in files:
lines = f["summary"]["lines"]
acc = crates.setdefault(crate_label(f["filename"], root), [0, 0])
acc[0] += lines["covered"]
acc[1] += lines["count"]
rows = sorted(
crates.items(),
key=lambda kv: (100.0 * kv[1][0] / kv[1][1]) if kv[1][1] else 101.0,