Compare commits

..

8 Commits

Author SHA1 Message Date
houseme 6c99df80f2 Merge branch 'main' into overtrue/report-security-coverage-ratchet 2026-08-23 12:32:30 +08:00
Zhengchao An 23a2c7d776 test(kms): stabilize Vault failover validation (#6385)
* test(kms): bound Vault failover progress wait

* ci(nightly): honor manual dispatch ref

* test(kms): preserve Vault worker failures

* test(kms): validate Vault circuit recovery
2026-08-23 12:32:11 +08:00
overtrue 2786c6813a ci(coverage): cancel superseded PR runs 2026-08-23 04:31:40 +08:00
overtrue d3fe16442e fix(ci): reject malformed coverage counts 2026-08-23 04:09:21 +08:00
overtrue 60a05220d2 fix(ci): reject incomplete coverage baselines 2026-08-23 03:19:23 +08:00
overtrue d37746fcc3 chore: merge main into coverage ratchet 2026-08-23 02:41:10 +08:00
overtrue d614fdbe34 ci(coverage): preserve report publication margin 2026-08-23 00:34:18 +08:00
overtrue 09de44df08 ci(coverage): add security ratchet calibration 2026-08-22 20:43:36 +08:00
14 changed files with 522 additions and 973 deletions
+20
View File
@@ -0,0 +1,20 @@
# 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,6 +36,7 @@ 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
+28 -12
View File
@@ -12,14 +12,12 @@
# See the License for the specific language governing permissions and
# limitations under the License.
# Weekly workspace line-coverage baseline (backlog#1153 infra-5).
# Workspace line-coverage baseline and security-crate calibration
# (backlog#1153 infra-5/infra-6).
#
# 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).
# 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.
#
# Measurement scope matches the PR test gate (ci.yml "Run tests"):
# `--workspace --exclude e2e_test` with the `ci` nextest profile. Doctests are
@@ -31,6 +29,17 @@
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),
@@ -39,6 +48,10 @@ 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:
@@ -46,12 +59,13 @@ permissions:
jobs:
coverage:
name: Workspace coverage (weekly)
name: Workspace line coverage
runs-on: sm-standard-4
# The instrumented build cannot reuse the regular CI cache (different
# 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
# 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
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
# Match the PR gate's nextest semantics (ci.yml runs `--profile ci`):
@@ -91,7 +105,9 @@ 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"
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"
- name: Upload coverage artifact
if: always()
+3 -6
View File
@@ -39,11 +39,10 @@ jobs:
env:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -89,11 +88,10 @@ jobs:
# either casing.
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
@@ -178,11 +176,10 @@ jobs:
FORCE_JAVASCRIPT_ACTIONS_TO_NODE24: "true"
NO_PROXY: 127.0.0.1,localhost
steps:
- name: Checkout main branch
- name: Checkout repository
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
with:
persist-credentials: false
ref: main
- name: Setup Rust environment
uses: ./.github/actions/setup
+1 -632
View File
@@ -21,7 +21,6 @@ use arc_swap::ArcSwapOption;
use rmp::Marker;
use serde::{Deserialize, Serialize};
use std::cmp::Ordering;
use std::collections::{HashMap, HashSet};
use std::str::from_utf8;
use std::{
fmt::Debug,
@@ -38,13 +37,8 @@ use tokio::io::{AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt};
use tokio::spawn;
use tokio::sync::Mutex;
use tracing::{debug, warn};
use uuid::Uuid;
const SLASH_SEPARATOR: &str = "/";
pub const MAX_META_CACHE_HEAL_CANDIDATES: usize = 1024;
/// Keep truncation continuations bounded while still giving the scanner a
/// safe object-level retry for versions that did not fit in the candidate set.
pub const MAX_META_CACHE_HEAL_TRUNCATED_OBJECTS: usize = 64;
#[derive(Clone, Debug, Default)]
pub struct MetadataResolutionParams {
@@ -72,50 +66,6 @@ pub struct MetaCacheEntry {
pub reusable: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub enum MetaCacheHealCandidateKind {
Object,
DeleteMarker,
UnversionedObject,
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct MetaCacheHealCandidate {
pub object: String,
pub version_id: Option<Uuid>,
pub kind: MetaCacheHealCandidateKind,
/// Number of raw disk entries that carried this validated version.
pub replica_count: usize,
}
impl MetaCacheHealCandidate {
pub fn validated_version(&self) -> Option<Uuid> {
match self.kind {
MetaCacheHealCandidateKind::Object | MetaCacheHealCandidateKind::DeleteMarker => self.version_id,
MetaCacheHealCandidateKind::UnversionedObject => None,
}
}
pub fn is_unversioned(&self) -> bool {
self.kind == MetaCacheHealCandidateKind::UnversionedObject
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct MetaCacheHealDiscovery {
pub candidates: Vec<MetaCacheHealCandidate>,
pub unverified_count: usize,
pub truncated: bool,
/// Object names whose validated version set exceeded the candidate cap.
/// The scanner retries these names without a version and with destructive
/// healing disabled; this is an explicit bounded continuation, not a
/// version claim.
pub truncated_objects: Vec<String>,
/// Validated candidates beyond the main cap, retained with exact version
/// identities so callers never fall back to a latest-version request.
pub truncated_candidates: Vec<MetaCacheHealCandidate>,
}
impl MetaCacheEntry {
pub fn marshal_msg(&self) -> Result<Vec<u8>> {
let mut wr = Vec::new();
@@ -420,185 +370,6 @@ impl MetaCacheEntries {
})
}
/// Discover validated object/delete-marker versions and safe unversioned
/// inspection candidates in the raw entries without applying read quorum.
/// This is intentionally separate from [`Self::resolve`]: a sub-quorum
/// version is a valid heal target even though it must not participate in
/// normal reads or writes.
///
/// The validated list is bounded and deduplicated by object, version id,
/// and metadata kind; each candidate retains the number of raw disk
/// entries that carried it so callers can classify sub-quorum versions.
/// Entries whose xl.meta cannot be decoded are counted separately for
/// discovery accounting; they never become versionless destructive heal
/// requests and do not consume the validated quota. An
/// [`MetaCacheHealCandidateKind::UnversionedObject`] is always consumed by
/// a non-destructive scanner request.
pub fn discover_heal_candidates(&self, bucket: &str, max_candidates: usize) -> MetaCacheHealDiscovery {
let limit = max_candidates.min(MAX_META_CACHE_HEAL_CANDIDATES);
if limit == 0 || bucket.is_empty() {
return MetaCacheHealDiscovery::default();
}
let mut discovery = MetaCacheHealDiscovery {
candidates: Vec::<MetaCacheHealCandidate>::with_capacity(limit.min(self.0.len())),
unverified_count: 0,
truncated: false,
truncated_objects: Vec::with_capacity(MAX_META_CACHE_HEAL_TRUNCATED_OBJECTS.min(limit)),
truncated_candidates: Vec::new(),
};
let mut seen: HashMap<(String, Option<Uuid>, MetaCacheHealCandidateKind), usize> =
HashMap::with_capacity(limit.min(self.0.len()));
for entry in self.0.iter().flatten() {
if !valid_heal_candidate_name(bucket, entry) {
continue;
}
let meta = match FileMeta::load(&entry.metadata) {
Ok(meta) => meta,
Err(_) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
let mut entry_seen = HashSet::new();
for shallow in meta.versions {
let version = match shallow.parse_version_meta() {
Ok(version) if version.valid() => version,
Ok(_) | Err(_) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
if version.free_version() {
continue;
}
let payload_header = version.header();
if normalize_version_id(shallow.header.version_id) != normalize_version_id(payload_header.version_id)
|| shallow.header.version_type != payload_header.version_type
{
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
let (kind, version_id) = match version.version_type {
VersionType::Object
if version.object.is_some() && version.delete_marker.is_none() && version.legacy_object.is_none() =>
{
match version.object.as_ref().and_then(|object| object.version_id) {
Some(id) if !id.is_nil() => (MetaCacheHealCandidateKind::Object, Some(id)),
Some(_) | None => (MetaCacheHealCandidateKind::UnversionedObject, None),
}
}
VersionType::Delete
if version.delete_marker.is_some() && version.object.is_none() && version.legacy_object.is_none() =>
{
let Some(id) = version.delete_marker.as_ref().and_then(|marker| marker.version_id) else {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
};
if id.is_nil() {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
(MetaCacheHealCandidateKind::DeleteMarker, Some(id))
}
VersionType::Legacy
if version.legacy_object.is_some() && version.object.is_none() && version.delete_marker.is_none() =>
{
let Some(legacy) = version.legacy_object.as_ref() else {
continue;
};
if legacy.version_id.is_empty() {
(MetaCacheHealCandidateKind::UnversionedObject, None)
} else {
let Ok(id) = Uuid::parse_str(&legacy.version_id) else {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
};
if id.is_nil() {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
(MetaCacheHealCandidateKind::Object, Some(id))
}
}
_ => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
if normalize_version_id(payload_header.version_id) != version_id {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
// `all_parts=true` is the trust-boundary check for versioned
// candidates. A null/legacy object may still need the old
// non-destructive inspection fallback when its part arrays
// are parseable but incomplete; never use that fallback for
// a candidate carrying a real version id.
let file_info = match version.clone().into_fileinfo(bucket, &entry.name, true) {
Ok(file_info) => file_info,
Err(_) if version_id.is_none() && matches!(kind, MetaCacheHealCandidateKind::UnversionedObject) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
match version.into_fileinfo(bucket, &entry.name, false) {
Ok(file_info) => file_info,
Err(_) => continue,
}
}
Err(_) => {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
};
if file_info.volume != bucket || file_info.name != entry.name {
discovery.unverified_count = discovery.unverified_count.saturating_add(1);
continue;
}
let candidate = MetaCacheHealCandidate {
object: entry.name.clone(),
version_id,
kind,
replica_count: 1,
};
let key = (candidate.object.clone(), candidate.version_id, candidate.kind.clone());
if entry_seen.contains(&key) {
continue;
}
if let Some(index) = seen.get(&key).copied() {
entry_seen.insert(key);
discovery.candidates[index].replica_count = discovery.candidates[index].replica_count.saturating_add(1);
} else if discovery.candidates.len() >= limit {
// Keep the validated candidate list bounded, but retain a
// bounded object-level continuation so the scanner cannot
// silently lose every version of a busy object.
discovery.truncated = true;
if discovery.truncated_objects.len() < MAX_META_CACHE_HEAL_TRUNCATED_OBJECTS
&& !discovery.truncated_objects.iter().any(|object| object == &candidate.object)
{
discovery.truncated_objects.push(candidate.object.clone());
}
if discovery.truncated_objects.iter().any(|object| object == &candidate.object) {
discovery.truncated_candidates.push(candidate);
}
continue;
} else {
entry_seen.insert(key.clone());
seen.insert(key, discovery.candidates.len());
discovery.candidates.push(candidate);
}
}
}
discovery
}
fn resolve_inner(&self, mut params: MetadataResolutionParams, enforce_write_quorum: bool) -> Option<MetaCacheEntry> {
if self.0.is_empty() {
debug!(
@@ -775,33 +546,6 @@ impl MetaCacheEntries {
}
}
fn valid_heal_candidate_name(bucket: &str, entry: &MetaCacheEntry) -> bool {
if bucket.is_empty()
|| entry.name.is_empty()
|| entry.is_dir()
|| (cfg!(windows) && entry.name.contains('\\'))
|| entry.name.chars().any(char::is_control)
{
return false;
}
// Validate raw key components without normalizing them. The scanner maps
// accepted keys to filesystem paths later, so dot components and empty
// internal components must be rejected before that boundary. A final
// empty component is retained for valid keys ending in '/'.
let mut components = entry.name.split('/').peekable();
while let Some(component) = components.next() {
if component == "." || component == ".." || (component.is_empty() && components.peek().is_some()) {
return false;
}
}
true
}
fn normalize_version_id(version_id: Option<Uuid>) -> Option<Uuid> {
version_id.filter(|id| !id.is_nil())
}
#[derive(Debug, Default)]
pub struct MetaCacheEntriesSortedResult {
pub entries: Option<MetaCacheEntriesSorted>,
@@ -1247,7 +991,7 @@ impl<T: Clone + Debug + Send + Sync + 'static> Cache<T> {
mod tests {
use super::*;
use crate::test_data::create_real_xlmeta;
use crate::{FileMetaVersion, MetaDeleteMarker, MetaObjectV1, MetaObjectV1Erasure, MetaObjectV1Stat, TRANSITION_COMPLETE};
use crate::{FileMetaVersion, MetaDeleteMarker, TRANSITION_COMPLETE};
use std::collections::HashMap;
use std::io::Cursor;
use std::sync::{
@@ -1848,381 +1592,6 @@ mod tests {
);
}
#[test]
fn discover_heal_candidates_keeps_sub_quorum_versions_and_deduplicates() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let entries = MetaCacheEntries(vec![
Some(metacache_entry_single_version(1, now, "one")),
Some(metacache_entry_single_version(2, now, "two")),
Some(metacache_entry_single_version(2, now, "two")),
Some(metacache_entry_single_version(3, now, "three")),
]);
let discovery = entries.discover_heal_candidates("bucket", 16);
let ids: std::collections::HashSet<Uuid> = discovery
.candidates
.iter()
.filter_map(|candidate| candidate.version_id)
.collect();
assert_eq!(
ids,
[Uuid::from_u128(1), Uuid::from_u128(2), Uuid::from_u128(3)]
.into_iter()
.collect()
);
assert_eq!(discovery.candidates.len(), 3, "duplicate tied versions must be emitted once");
assert_eq!(
discovery
.candidates
.iter()
.find(|candidate| candidate.version_id == Some(Uuid::from_u128(2)))
.expect("duplicate version should be discovered")
.replica_count,
2
);
}
#[test]
fn discover_heal_candidates_does_not_count_duplicate_versions_within_one_entry() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut meta = FileMeta::load(&metacache_entry_single_version(1, now, "duplicate").metadata)
.expect("duplicate fixture should decode");
meta.versions.push(meta.versions[0].clone());
let entry = MetaCacheEntry {
name: "object".to_string(),
metadata: meta.marshal_msg().expect("duplicate metadata should marshal"),
cached: Some(meta),
reusable: false,
};
let discovery = MetaCacheEntries(vec![Some(entry)]).discover_heal_candidates("bucket", 16);
let candidate = discovery
.candidates
.iter()
.find(|candidate| candidate.version_id == Some(Uuid::from_u128(1)))
.expect("duplicate fixture should be discovered");
assert_eq!(candidate.replica_count, 1);
}
#[test]
fn discover_heal_candidates_covers_divergent_quorum_boundaries_n2_n4_n6() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
for (disk_count, quorum) in [(2usize, 1usize), (4, 2), (6, 3)] {
let target_id = Uuid::from_u128(0x1000 + disk_count as u128);
for target_replicas in [quorum.saturating_sub(1), quorum, quorum + 1] {
let entries = (0..disk_count)
.map(|disk| {
let version_id = if disk < target_replicas {
target_id
} else {
Uuid::from_u128(0x2000 + disk as u128)
};
Some(metacache_entry_single_version(version_id.as_u128(), now, "divergent"))
})
.collect();
let discovery = MetaCacheEntries(entries).discover_heal_candidates("bucket", 32);
let target = discovery
.candidates
.iter()
.find(|candidate| candidate.version_id == Some(target_id));
assert_eq!(target.is_some(), target_replicas > 0, "N={disk_count}, replicas={target_replicas}");
if let Some(target) = target {
assert_eq!(target.replica_count, target_replicas);
}
}
}
}
#[test]
fn discover_heal_candidates_separates_delete_markers_and_preserves_unversioned_objects() {
let mut marker_meta = FileMeta::new();
marker_meta
.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(Uuid::from_u128(99)),
deleted: true,
mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")),
..Default::default()
})
.expect("delete marker should be added");
let marker = MetaCacheEntry {
name: "object".to_string(),
metadata: marker_meta.marshal_msg().expect("delete marker metadata should marshal"),
cached: Some(marker_meta),
reusable: false,
};
let unversioned_entry = metacache_entry_with_mod_time(
OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"),
"unversioned",
);
let discovery = MetaCacheEntries(vec![Some(marker), Some(unversioned_entry)]).discover_heal_candidates("bucket", 16);
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::DeleteMarker && candidate.version_id == Some(Uuid::from_u128(99))
}));
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none()
}));
}
#[test]
fn discover_heal_candidates_rejects_delete_markers_without_ids() {
let mut marker_meta = FileMeta::new();
marker_meta
.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
deleted: true,
mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")),
..Default::default()
})
.expect("nil delete marker should be added");
let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry {
name: "object".to_string(),
metadata: marker_meta.marshal_msg().expect("nil marker metadata should marshal"),
cached: Some(marker_meta),
reusable: false,
})])
.discover_heal_candidates("bucket", 16);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.kind == MetaCacheHealCandidateKind::DeleteMarker)
);
assert!(discovery.unverified_count >= 1);
}
#[test]
fn discover_heal_candidates_skips_free_versions() {
let object_id = Uuid::from_u128(100);
let free_id = Uuid::from_u128(101);
let mut meta = FileMeta::new();
meta.add_version(FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(object_id),
transition_status: TRANSITION_COMPLETE.to_string(),
transitioned_objname: "remote/object".to_string(),
transition_version_id: Some(Uuid::from_u128(102)),
transition_tier: "WARM".to_string(),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
})
.expect("transitioned object should be added");
let mut delete = FileInfo {
volume: "bucket".to_string(),
name: "object".to_string(),
version_id: Some(object_id),
mod_time: Some(OffsetDateTime::now_utc()),
..Default::default()
};
delete.set_tier_free_version_id(&free_id.to_string());
meta.delete_version(&delete).expect("free version should be persisted");
let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry {
name: "object".to_string(),
metadata: meta.marshal_msg().expect("free version metadata should marshal"),
cached: Some(meta),
reusable: false,
})])
.discover_heal_candidates("bucket", 16);
assert!(discovery.candidates.is_empty());
}
#[test]
fn discover_heal_candidates_preserves_unversioned_legacy_object() {
let legacy = MetaObjectV1 {
version: "1.0.1".to_string(),
format: "xl".to_string(),
stat: MetaObjectV1Stat {
size: 1,
mod_time: Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp")),
name: "object".to_string(),
..Default::default()
},
erasure: MetaObjectV1Erasure {
data_blocks: 4,
parity_blocks: 2,
index: 1,
distribution: vec![1, 2, 3, 4, 5, 6],
..Default::default()
},
..Default::default()
};
let version = FileMetaVersion {
version_type: VersionType::Legacy,
legacy_object: Some(legacy),
..Default::default()
};
let mut meta = FileMeta::new();
meta.versions
.push(FileMetaShallowVersion::try_from(version).expect("legacy metadata should marshal"));
let discovery = MetaCacheEntries(vec![Some(MetaCacheEntry {
name: "object".to_string(),
metadata: meta.marshal_msg().expect("legacy metadata should marshal"),
cached: Some(meta),
reusable: false,
})])
.discover_heal_candidates("bucket", 16);
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none()
}));
}
#[test]
fn discover_heal_candidates_rejects_nil_and_malformed_metadata_and_is_bounded() {
let now = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut nil = metacache_entry_single_version(1, now, "nil");
let mut nil_meta = FileMeta::load(&nil.metadata).expect("nil fixture should decode");
let mut nil_version = nil_meta.versions[0]
.parse_version_meta()
.expect("nil fixture version should decode");
nil_version.object.as_mut().expect("object fixture").version_id = Some(Uuid::nil());
nil_meta.versions[0] = FileMetaShallowVersion::try_from(nil_version).expect("nil fixture should marshal");
nil.metadata = nil_meta.marshal_msg().expect("nil fixture metadata should marshal");
let mut mismatched = metacache_entry_single_version(2, now, "mismatched");
let mut mismatched_meta = FileMeta::load(&mismatched.metadata).expect("mismatched fixture should decode");
mismatched_meta.versions[0].header.version_id = Some(Uuid::from_u128(200));
mismatched.metadata = mismatched_meta.marshal_msg().expect("mismatched metadata should marshal");
let mut short_parts = metacache_entry_single_version(3, now, "short-parts");
let mut short_parts_meta = FileMeta::load(&short_parts.metadata).expect("short-parts fixture should decode");
let mut short_parts_version = short_parts_meta.versions[0]
.parse_version_meta()
.expect("short-parts fixture version should decode");
let object = short_parts_version.object.as_mut().expect("object fixture");
object.part_numbers = vec![1];
object.part_actual_sizes = vec![1];
object.part_sizes.clear();
short_parts_meta.versions[0] = FileMetaShallowVersion::try_from(short_parts_version).expect("short-parts should marshal");
short_parts.metadata = short_parts_meta.marshal_msg().expect("short-parts metadata should marshal");
let mut short_unversioned = metacache_entry_with_mod_time(now, "short-unversioned");
let mut short_unversioned_meta =
FileMeta::load(&short_unversioned.metadata).expect("short-unversioned fixture should decode");
let mut short_unversioned_version = short_unversioned_meta.versions[0]
.parse_version_meta()
.expect("short-unversioned version should decode");
let unversioned_object = short_unversioned_version.object.as_mut().expect("unversioned object fixture");
unversioned_object.part_numbers = vec![1];
unversioned_object.part_actual_sizes = vec![1];
unversioned_object.part_sizes.clear();
short_unversioned_meta.versions[0] =
FileMetaShallowVersion::try_from(short_unversioned_version).expect("short-unversioned should marshal");
short_unversioned.metadata = short_unversioned_meta
.marshal_msg()
.expect("short-unversioned metadata should marshal");
let mut malformed = nil.clone();
malformed.name = "malformed".to_string();
malformed.metadata = vec![1, 2, 3];
let entries = MetaCacheEntries(
std::iter::once(Some(nil))
.chain(std::iter::once(Some(mismatched)))
.chain(std::iter::once(Some(short_parts)))
.chain(std::iter::once(Some(short_unversioned)))
.chain(std::iter::once(Some(malformed)))
.chain((0..32).map(|id| Some(metacache_entry_single_version(id + 10, now, "bounded"))))
.collect(),
);
let discovery = entries.discover_heal_candidates("bucket", 5);
assert!(discovery.candidates.len() <= 5);
assert!(discovery.truncated, "bounded discovery must expose dropped candidates");
assert!(
discovery
.truncated_candidates
.iter()
.all(|candidate| candidate.version_id.is_some()),
"overflow candidates must retain exact version identities"
);
assert!(
discovery.truncated_objects.iter().any(|object| object == "object"),
"bounded discovery must expose an object-level safe continuation"
);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.version_id == Some(Uuid::nil()))
);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.version_id == Some(Uuid::from_u128(2)))
);
assert!(
!discovery
.candidates
.iter()
.any(|candidate| candidate.version_id == Some(Uuid::from_u128(3)))
);
assert!(discovery.candidates.iter().any(|candidate| {
candidate.kind == MetaCacheHealCandidateKind::UnversionedObject && candidate.version_id.is_none()
}));
assert!(
discovery.unverified_count >= 1,
"malformed and rejected metadata must remain observable during discovery"
);
for invalid_name in [
"../object",
"./object",
"object/../other",
"object//name",
"object\u{0001}name",
"object\0name",
] {
let mut entry = metacache_entry_single_version(400, now, invalid_name);
entry.name = invalid_name.to_string();
let discovery = MetaCacheEntries(vec![Some(entry)]).discover_heal_candidates("bucket", 5);
assert!(
discovery.candidates.is_empty(),
"invalid key should not become a heal candidate: {invalid_name:?}"
);
}
#[cfg(windows)]
{
let mut entry = metacache_entry_single_version(400, now, "object\\name");
entry.name = "object\\name".to_string();
assert!(
MetaCacheEntries(vec![Some(entry)])
.discover_heal_candidates("bucket", 5)
.candidates
.is_empty(),
"backslash is a path separator on Windows"
);
}
#[cfg(not(windows))]
{
let mut entry = metacache_entry_single_version(400, now, "object\\name");
entry.name = "object\\name".to_string();
assert_eq!(
MetaCacheEntries(vec![Some(entry)])
.discover_heal_candidates("bucket", 5)
.candidates
.len(),
1,
"backslash is object-key data on Unix"
);
}
for valid_name in ["trailing/", "prefix/object"] {
let mut entry = metacache_entry_single_version(401, now, valid_name);
entry.name = valid_name.to_string();
let discovery = MetaCacheEntries(vec![Some(entry)]).discover_heal_candidates("bucket", 5);
assert_eq!(discovery.candidates.len(), 1, "raw S3 key should remain opaque: {valid_name:?}");
}
}
#[test]
fn resolve_rejects_partial_latest_and_returns_committed_previous_metadata() {
let old_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
+86 -25
View File
@@ -16,14 +16,14 @@
//!
//! `scripts/test/vault_ha_kms_live.sh` owns the official Vault containers and
//! kills the active node while this test continuously decrypts through a
//! surviving standby. KV2 and Transit requests must remain successful, use a
//! bounded number of attempts, and leave the circuit and in-flight gauges at
//! zero after a new leader is elected.
//! surviving standby. KV2 and Transit must recover after the bounded circuit
//! interval, use a bounded number of attempts, and leave the circuit and
//! in-flight gauges at zero after a new leader is elected.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use metrics_util::MetricKind;
@@ -43,6 +43,11 @@ const OPERATION_ATTEMPTS: &str = "rustfs_kms_backend_operation_attempts";
const IN_FLIGHT: &str = "rustfs_kms_backend_in_flight";
const CIRCUIT_OPEN: &str = "rustfs_kms_backend_circuit_open";
const MAX_ATTEMPTS: u32 = 10;
const ATTEMPT_TIMEOUT: Duration = Duration::from_secs(2);
const HEALTHY_PROGRESS_TIMEOUT: Duration = Duration::from_secs(20);
// The circuit remains open for 30s after five failed attempts.
const POST_FAILOVER_PROGRESS_TIMEOUT: Duration = Duration::from_secs(35);
const FAILOVER_ERROR_POLL_INTERVAL: Duration = Duration::from_millis(100);
type MetricEntry = (
metrics_util::CompositeKey,
@@ -64,7 +69,7 @@ fn config(backend: KmsBackend, backend_config: BackendConfig) -> KmsConfig {
backend,
backend_config,
allow_insecure_dev_defaults: true,
timeout: Duration::from_secs(2),
timeout: ATTEMPT_TIMEOUT,
retry_attempts: MAX_ATTEMPTS,
enable_cache: false,
..KmsConfig::default()
@@ -164,14 +169,31 @@ fn retryable_failures(snapshot: &[MetricEntry], operation: &str) -> u64 {
.sum()
}
async fn wait_for_count(counter: &AtomicU64, minimum: u64, description: &str) {
tokio::time::timeout(Duration::from_secs(20), async {
async fn wait_for_count(
counter: &AtomicU64,
failure: &Mutex<Option<String>>,
minimum: u64,
description: &str,
timeout: Duration,
) {
tokio::time::timeout(timeout, async {
while counter.load(Ordering::SeqCst) < minimum {
if let Some(error) = failure.lock().expect("decrypt failure lock poisoned").as_ref() {
panic!(
"{description} worker failed after {} successful decrypts: {error}",
counter.load(Ordering::SeqCst)
);
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await
.unwrap_or_else(|_| panic!("timed out waiting for {description}"));
.unwrap_or_else(|_| {
panic!(
"timed out after {timeout:?} waiting for {description}: completed {}, expected {minimum}",
counter.load(Ordering::SeqCst)
)
});
}
async fn wait_for_file(path: &Path, description: &str) {
@@ -189,7 +211,8 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
request: DecryptRequest,
expected: Vec<u8>,
completed: Arc<AtomicU64>,
failed: Arc<AtomicBool>,
allow_failover_errors: Arc<AtomicBool>,
failure: Arc<Mutex<Option<String>>>,
stop: CancellationToken,
) {
while !stop.is_cancelled() {
@@ -197,8 +220,18 @@ async fn decrypt_loop<B: KmsBackendTrait + Send + Sync + 'static>(
Ok(response) if response.plaintext == expected => {
completed.fetch_add(1, Ordering::SeqCst);
}
Ok(_) | Err(_) => {
failed.store(true, Ordering::SeqCst);
Ok(_) => {
*failure.lock().expect("decrypt failure lock poisoned") =
Some("decrypt returned unexpected plaintext".to_string());
return;
}
Err(rustfs_kms::KmsError::BackendError { .. } | rustfs_kms::KmsError::OperationTimedOut { .. })
if allow_failover_errors.load(Ordering::SeqCst) =>
{
tokio::time::sleep(FAILOVER_ERROR_POLL_INTERVAL).await;
}
Err(error) => {
*failure.lock().expect("decrypt failure lock poisoned") = Some(error.to_string());
return;
}
}
@@ -296,7 +329,9 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
);
let stop = CancellationToken::new();
let failed = Arc::new(AtomicBool::new(false));
let allow_failover_errors = Arc::new(AtomicBool::new(false));
let kv2_failure = Arc::new(Mutex::new(None));
let transit_failure = Arc::new(Mutex::new(None));
let kv2_completed = Arc::new(AtomicU64::new(0));
let transit_completed = Arc::new(AtomicU64::new(0));
let kv2_worker = tokio::spawn(decrypt_loop(
@@ -304,7 +339,8 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
kv2_request,
kv2_data_key.plaintext_key,
Arc::clone(&kv2_completed),
Arc::clone(&failed),
Arc::clone(&allow_failover_errors),
Arc::clone(&kv2_failure),
stop.clone(),
));
let transit_worker = tokio::spawn(decrypt_loop(
@@ -312,12 +348,21 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
transit_request,
transit_data_key.plaintext_key,
Arc::clone(&transit_completed),
Arc::clone(&failed),
Arc::clone(&allow_failover_errors),
Arc::clone(&transit_failure),
stop.clone(),
));
wait_for_count(&kv2_completed, 2, "two healthy KV2 decrypts").await;
wait_for_count(&transit_completed, 2, "two healthy Transit decrypts").await;
wait_for_count(&kv2_completed, &kv2_failure, 2, "two healthy KV2 decrypts", HEALTHY_PROGRESS_TIMEOUT).await;
wait_for_count(
&transit_completed,
&transit_failure,
2,
"two healthy Transit decrypts",
HEALTHY_PROGRESS_TIMEOUT,
)
.await;
allow_failover_errors.store(true, Ordering::SeqCst);
std::fs::write(&marker, b"ready").expect("publish failover readiness marker");
wait_for_file(&elected, "the replacement Vault leader").await;
@@ -326,18 +371,39 @@ async fn exercise_failover(snapshotter: &Snapshotter) {
let kv2_after_election = kv2_completed.load(Ordering::SeqCst) + 2;
let transit_after_election = transit_completed.load(Ordering::SeqCst) + 2;
wait_for_count(&kv2_completed, kv2_after_election, "post-failover KV2 decrypts").await;
wait_for_count(&transit_completed, transit_after_election, "post-failover Transit decrypts").await;
wait_for_count(
&kv2_completed,
&kv2_failure,
kv2_after_election,
"post-failover KV2 decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
wait_for_count(
&transit_completed,
&transit_failure,
transit_after_election,
"post-failover Transit decrypts",
POST_FAILOVER_PROGRESS_TIMEOUT,
)
.await;
stop.cancel();
kv2_worker.await.expect("KV2 decrypt worker must join");
transit_worker.await.expect("Transit decrypt worker must join");
assert!(!failed.load(Ordering::SeqCst), "no decrypt may fail or return different plaintext");
assert!(
kv2_failure.lock().expect("KV2 failure lock poisoned").is_none(),
"no KV2 decrypt may fail or return different plaintext"
);
assert!(
transit_failure.lock().expect("Transit failure lock poisoned").is_none(),
"no Transit decrypt may fail or return different plaintext"
);
}
#[test]
#[ignore = "requires a real three-node Vault Raft cluster; run scripts/test/vault_ha_kms_live.sh"]
fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
fn vault_raft_leader_failure_recovers_kv2_and_transit_decrypts() {
let recorder = DebuggingRecorder::new();
let snapshotter = recorder.snapshotter();
metrics::with_local_recorder(&recorder, || {
@@ -349,11 +415,6 @@ fn vault_raft_leader_failure_preserves_kv2_and_transit_decrypts() {
});
let snapshot = snapshotter.snapshot().into_vec();
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "circuit_open")]),
0,
"a bounded leader election must not open the circuit"
);
assert_eq!(
counter_value(&snapshot, OPERATIONS_TOTAL, &[("outcome", "budget_exhausted")]),
0,
+67 -116
View File
@@ -43,10 +43,7 @@ use rustfs_common::metrics::{
UpdateCurrentPathFn, current_path_updater, global_metrics,
};
use rustfs_common::trace_bus::{TraceEvent, TraceFunc, TraceKind, trace_emit, trace_subscriber_count};
use rustfs_filemeta::{
MAX_META_CACHE_HEAL_CANDIDATES, MAX_META_CACHE_HEAL_TRUNCATED_OBJECTS, MetaCacheEntries, MetaCacheEntry,
MetaCacheHealCandidateKind,
};
use rustfs_filemeta::{MetaCacheEntries, MetaCacheEntry, MetadataResolutionParams};
use rustfs_utils::path::{SLASH_SEPARATOR, path_join_buf};
use s3s::dto::{BucketLifecycleConfiguration, ObjectLockConfiguration, VersioningConfiguration};
use time::OffsetDateTime;
@@ -99,11 +96,6 @@ const METRIC_SCANNER_EXCESS_OBJECT_VERSION_SIZE_TOTAL: &str = "rustfs_scanner_ex
const METRIC_SCANNER_EXCESS_FOLDERS_TOTAL: &str = "rustfs_scanner_excess_folders_total";
const METRIC_SCANNER_PENDING_HEAL_PRUNE_TOTAL: &str = "rustfs_scanner_pending_heal_prune_total";
const METRIC_SCANNER_PENDING_HEAL_MALFORMED_TOTAL: &str = "rustfs_scanner_pending_heal_malformed_total";
const METRIC_SCANNER_HEAL_DISCOVERY_CANDIDATES_TOTAL: &str = "rustfs_scanner_heal_discovery_candidates_total";
const METRIC_SCANNER_HEAL_DISCOVERY_SUB_QUORUM_TOTAL: &str = "rustfs_scanner_heal_discovery_sub_quorum_total";
const METRIC_SCANNER_HEAL_DISCOVERY_UNVERIFIED_TOTAL: &str = "rustfs_scanner_heal_discovery_unverified_total";
const METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL: &str = "rustfs_scanner_heal_discovery_queued_total";
const METRIC_SCANNER_HEAL_DISCOVERY_TRUNCATED_TOTAL: &str = "rustfs_scanner_heal_discovery_truncated_total";
const MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET: usize = 128;
// --- scanner excess alerts as S3 notification events (rustfs/backlog#1868) --
@@ -891,7 +883,7 @@ impl FolderScanner {
object: Option<String>,
version_id: Option<String>,
request: HealChannelRequest,
) -> Result<HealAdmissionResult, ScannerError> {
) -> Result<(), ScannerError> {
let candidate_type = pending_scanner_heal_candidate_type(kind);
let priority = request.priority;
let scan_mode = request.scan_mode.unwrap_or(self.scan_mode);
@@ -919,7 +911,7 @@ impl FolderScanner {
error = %err,
"Scanner deferred heal request after channel error"
);
return Ok(HealAdmissionResult::Full);
return Ok(());
}
};
self.update_pending_scanner_heal_after_admission(
@@ -931,7 +923,7 @@ impl FolderScanner {
result,
);
if result.is_admitted() {
return Ok(result);
return Ok(());
}
record_high_priority_heal_escalation(candidate_type, priority, result);
@@ -952,7 +944,7 @@ impl FolderScanner {
state = "high_priority_not_admitted",
"Scanner high-priority heal admission failed"
);
Ok(result)
Ok(())
}
pub fn set_heal_object_select(&mut self, prob: u32) {
@@ -1744,7 +1736,14 @@ impl FolderScanner {
break;
}
let mut previous_bucket = String::new();
let mut resolver = MetadataResolutionParams {
dir_quorum: self.disks_quorum,
obj_quorum: self.disks_quorum,
bucket: "".to_string(),
strict: false,
..Default::default()
};
for name in abandoned_children {
if !self.should_heal().await {
break;
@@ -1752,7 +1751,7 @@ impl FolderScanner {
let (bucket, prefix) = path2_bucket_object(name.as_str());
if bucket != previous_bucket {
if bucket != resolver.bucket {
self.send_required_scanner_heal_request(
PendingScannerHealKind::Bucket,
bucket.clone(),
@@ -1761,9 +1760,10 @@ impl FolderScanner {
build_bucket_heal_request(bucket.clone(), HealChannelPriority::High),
)
.await?;
previous_bucket = bucket.clone();
}
resolver.bucket = bucket.clone();
let child_ctx = ctx.child_token();
let (agreed_tx, mut agreed_rx) = mpsc::channel::<String>(1);
@@ -1880,8 +1880,6 @@ impl FolderScanner {
let mut agreed_closed = false;
let mut partial_closed = false;
let mut finished_closed = false;
let mut seen_heal_candidates: HashSet<(String, Option<String>, MetaCacheHealCandidateKind)> = HashSet::new();
let mut seen_truncated_objects: HashSet<String> = HashSet::new();
loop {
if agreed_closed && partial_closed && finished_closed {
@@ -1906,112 +1904,65 @@ impl FolderScanner {
break;
}
let discovery = entries.discover_heal_candidates(&bucket, MAX_META_CACHE_HEAL_CANDIDATES);
counter!(METRIC_SCANNER_HEAL_DISCOVERY_CANDIDATES_TOTAL)
.increment(u64::try_from(discovery.candidates.len()).unwrap_or(u64::MAX));
counter!(METRIC_SCANNER_HEAL_DISCOVERY_SUB_QUORUM_TOTAL).increment(
u64::try_from(
discovery
.candidates
.iter()
.filter(|candidate| candidate.replica_count < disks_quorum)
.count(),
)
.unwrap_or(u64::MAX),
);
counter!(METRIC_SCANNER_HEAL_DISCOVERY_UNVERIFIED_TOTAL).increment(
u64::try_from(discovery.unverified_count).unwrap_or(u64::MAX),
);
if discovery.truncated {
counter!(METRIC_SCANNER_HEAL_DISCOVERY_TRUNCATED_TOTAL).increment(1);
let Some(entry) = resolve_object_heal_entry(&entries, resolver.clone()) else {
continue;
};
(self.update_current_path)(&entry.name).await;
if entry.is_dir() {
continue;
}
for candidate in discovery.candidates {
let sub_quorum_candidate = candidate.replica_count < disks_quorum;
let version_id = candidate.validated_version().map(|id| id.to_string());
let identity = (candidate.object.clone(), version_id.clone(), candidate.kind.clone());
if seen_heal_candidates.len() >= MAX_META_CACHE_HEAL_CANDIDATES
&& !seen_heal_candidates.contains(&identity)
{
continue;
}
if !seen_heal_candidates.insert(identity) {
continue;
}
let request = if candidate.is_unversioned() {
build_non_destructive_object_heal_request(
bucket.clone(),
candidate.object.clone(),
self.scan_mode,
HealChannelPriority::High,
)
} else {
build_object_heal_request(
bucket.clone(),
candidate.object.clone(),
version_id.clone(),
self.scan_mode,
HealChannelPriority::High,
)
};
(self.update_current_path)(&candidate.object).await;
let admission = self.send_required_scanner_heal_request(
PendingScannerHealKind::Object,
bucket.clone(),
Some(candidate.object.clone()),
version_id.clone(),
request,
)
.await?;
if admission.is_admitted() {
counter!(METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL).increment(1);
} else if sub_quorum_candidate {
self.mark_pending_scanner_heal_reason(
PendingScannerHealKind::Object,
&bucket,
Some(&candidate.object),
version_id.as_deref(),
"sub_quorum_metadata",
let fivs = match entry.file_info_versions(&bucket) {
Ok(fivs) => fivs,
Err(e) => {
error!(
target: "rustfs::scanner::folder",
event = EVENT_SCANNER_FOLDER_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_FOLDER,
bucket = %bucket,
entry = %entry.name,
state = "file_info_versions_failed",
error = %e,
"Scanner list_path_raw failed to resolve file versions"
);
}
found_objects = true;
}
// Candidates beyond the main cap remain exact
// version requests; never downgrade them to a
// latest-version (version_id=None) heal.
for candidate in discovery.truncated_candidates {
let version_id = candidate.validated_version().map(|id| id.to_string());
let identity = (candidate.object.clone(), version_id.clone(), candidate.kind.clone());
if seen_truncated_objects.len() >= MAX_META_CACHE_HEAL_TRUNCATED_OBJECTS
&& !seen_truncated_objects.contains(&candidate.object)
{
continue;
}
seen_truncated_objects.insert(candidate.object.clone());
if !seen_heal_candidates.insert(identity) {
continue;
}
let request = build_object_heal_request(
bucket.clone(),
candidate.object.clone(),
version_id.clone(),
self.scan_mode,
HealChannelPriority::High,
);
(self.update_current_path)(&candidate.object).await;
let admission = self
.send_required_scanner_heal_request(
self.send_required_scanner_heal_request(
PendingScannerHealKind::Object,
bucket.clone(),
Some(candidate.object.clone()),
version_id,
request,
Some(entry.name.clone()),
None,
build_object_heal_request(
bucket.clone(),
entry.name.clone(),
None,
self.scan_mode,
HealChannelPriority::High,
),
)
.await?;
if admission.is_admitted() {
counter!(METRIC_SCANNER_HEAL_DISCOVERY_QUEUED_TOTAL).increment(1);
found_objects = true;
continue;
}
};
for fiv in fivs.versions {
let version_id = fiv.version_id.and_then(|v| if v.is_nil() { None } else { Some(v.to_string()) });
self.send_required_scanner_heal_request(
PendingScannerHealKind::Object,
bucket.clone(),
Some(entry.name.clone()),
version_id.clone(),
build_object_heal_request(
bucket.clone(),
entry.name.clone(),
version_id,
self.scan_mode,
HealChannelPriority::High,
),
)
.await?;
found_objects = true;
}
@@ -13,8 +13,6 @@
// limitations under the License.
/// Per-object scan actions: ScannerItem, the get-size failure policy, and the heal/ILM admission helpers.
use super::*;
#[cfg(test)]
use rustfs_filemeta::MetadataResolutionParams;
/// Cached folder information for scanning
#[derive(Clone, Debug)]
@@ -90,22 +88,6 @@ pub(super) fn build_object_heal_request(
}
}
/// Build the versionless inspection request used when discovery cannot prove
/// a destructive version identity (for example an unversioned object or a
/// bounded candidate overflow). The explicit flag is the fail-closed safety
/// boundary; callers must not reconstruct it with the destructive default.
pub(super) fn build_non_destructive_object_heal_request(
bucket: String,
object: String,
scan_mode: HealScanMode,
priority: HealChannelPriority,
) -> HealChannelRequest {
let mut request = build_object_heal_request(bucket, object, None, scan_mode, priority);
request.remove_corrupted = Some(false);
request
}
#[cfg(test)]
pub(super) fn resolve_object_heal_entry(
entries: &MetaCacheEntries,
resolver: MetadataResolutionParams,
+7 -39
View File
@@ -105,29 +105,6 @@ impl FolderScanner {
}
}
/// Preserve the discovery reason when a candidate could not be admitted
/// immediately. The existing string field is intentionally reused so the
/// scanner's map-encoded cache schema stays backward compatible.
pub(super) fn mark_pending_scanner_heal_reason(
&mut self,
kind: PendingScannerHealKind,
bucket: &str,
object: Option<&str>,
version_id: Option<&str>,
reason: &str,
) {
if let Some(entry) = self
.new_cache
.info
.pending_heals
.iter_mut()
.find(|entry| pending_scanner_heal_matches(entry, kind, bucket, object, version_id))
{
entry.last_admission_reason = reason.to_string();
self.sync_pending_heals();
}
}
pub(super) fn prune_pending_scanner_heals(&mut self) {
let now = Self::now_secs();
let before_expiry = self.new_cache.info.pending_heals.len();
@@ -328,22 +305,13 @@ pub(super) fn build_pending_scanner_heal_request(entry: &PendingScannerHeal) ->
match entry.kind {
PendingScannerHealKind::Bucket => Some(build_bucket_heal_request(entry.bucket.clone(), HealChannelPriority::High)),
PendingScannerHealKind::Object => entry.object.as_ref().map(|object| {
if entry.version_id.is_none() {
build_non_destructive_object_heal_request(
entry.bucket.clone(),
object.clone(),
entry.scan_mode,
HealChannelPriority::High,
)
} else {
build_object_heal_request(
entry.bucket.clone(),
object.clone(),
entry.version_id.clone(),
entry.scan_mode,
HealChannelPriority::High,
)
}
build_object_heal_request(
entry.bucket.clone(),
object.clone(),
entry.version_id.clone(),
entry.scan_mode,
HealChannelPriority::High,
)
}),
}
}
+6 -106
View File
@@ -17,7 +17,7 @@ use crate::SCANNER_SLEEPER;
use super::*;
use crate::storage_api::VersionPurgeStatusType;
use crate::{DiskOption, Endpoint, STORAGE_FORMAT_FILE, TierStats, new_disk, storageclass};
use rustfs_filemeta::{FileInfo, FileMeta, MetadataResolutionParams};
use rustfs_filemeta::{FileInfo, FileMeta};
use std::io::Write;
#[cfg(unix)]
use std::os::unix::fs::{PermissionsExt, symlink};
@@ -982,21 +982,6 @@ fn test_build_object_heal_request_omits_nil_version_id() {
assert_eq!(request.recreate_missing, Some(false));
}
#[test]
fn test_build_non_destructive_object_heal_request_disables_removal() {
let request = build_non_destructive_object_heal_request(
"bucket".to_string(),
"path/to/object".to_string(),
HealScanMode::Deep,
HealChannelPriority::High,
);
assert_eq!(request.object_version_id, None);
assert_eq!(request.remove_corrupted, Some(false));
assert_eq!(request.recreate_missing, Some(false));
assert_eq!(request.source, HealRequestSource::Scanner);
}
#[test]
fn test_build_bucket_heal_request_disables_recreate_for_scanner() {
let request = build_bucket_heal_request("bucket".to_string(), HealChannelPriority::Low);
@@ -1136,42 +1121,6 @@ fn test_pending_heal_reconstructs_object_request_with_version() {
assert_eq!(request.source, HealRequestSource::Scanner);
}
#[test]
fn test_pending_heal_reconstructs_unversioned_request_without_removal() {
let pending = pending_heal(PendingScannerHealKind::Object, "bucket", Some("object"), None, 1, 1);
let request = build_pending_scanner_heal_request(&pending).expect("unversioned object request should rebuild");
assert!(request.object_version_id.is_none());
assert_eq!(request.remove_corrupted, Some(false));
assert_eq!(request.recreate_missing, Some(false));
}
#[tokio::test]
async fn test_pending_heal_reason_preserves_sub_quorum_discovery() {
let (mut scanner, temp_dir) = build_test_scanner().await;
let _guard = TestGuard::new(u64::MAX, usize::MAX, &mut scanner, temp_dir);
scanner.update_pending_scanner_heal_after_admission(
PendingScannerHealKind::Object,
"bucket",
Some("object"),
Some("version-a"),
HealScanMode::Deep,
HealAdmissionResult::Full,
);
scanner.mark_pending_scanner_heal_reason(
PendingScannerHealKind::Object,
"bucket",
Some("object"),
Some("version-a"),
"sub_quorum_metadata",
);
assert_eq!(scanner.new_cache.info.pending_heals.len(), 1);
assert_eq!(scanner.new_cache.info.pending_heals[0].last_admission_reason, "sub_quorum_metadata");
}
#[test]
fn test_pending_heal_retry_candidates_respect_cap_and_order() {
let pending: Vec<PendingScannerHeal> = (0..(MAX_PENDING_SCANNER_HEAL_RETRIES_PER_BUCKET + 2))
@@ -1374,20 +1323,6 @@ fn metadata_for_object(bucket: &str, object: &str) -> Vec<u8> {
meta.marshal_msg().expect("test metadata should marshal")
}
fn metadata_for_object_version(bucket: &str, object: &str, version_id: Option<Uuid>) -> Vec<u8> {
let mut file_info = FileInfo::new(object, 4, 2);
file_info.volume = bucket.to_string();
file_info.name = object.to_string();
file_info.version_id = version_id;
file_info.versioned = version_id.is_some();
file_info.mod_time = Some(OffsetDateTime::now_utc());
file_info.size = 1;
let mut meta = FileMeta::new();
meta.add_version(file_info).expect("test metadata version should be accepted");
meta.marshal_msg().expect("test metadata should marshal")
}
async fn write_test_object_metadata(root: &std::path::Path, bucket: &str, object: &str) {
write_test_object_metadata_bytes(root, bucket, object, &metadata_for_object(bucket, object)).await;
}
@@ -1791,21 +1726,12 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
let _guard = TestGuard::new(60, 100, &mut scanner, temp_dir.clone());
let heal_starts = Arc::new(AtomicUsize::new(0));
let heal_starts_clone = heal_starts.clone();
let healed_versions = Arc::new(Mutex::new(Vec::<Option<String>>::new()));
let healed_versions_clone = healed_versions.clone();
let mut heal_rx =
rustfs_common::heal_channel::init_heal_channel().expect("heal channel should initialize once for scanner tests");
let _heal_responder = tokio::spawn(async move {
while let Some(command) = heal_rx.recv().await {
if let rustfs_common::heal_channel::HealChannelCommand::Start {
request, response_tx, ..
} = command
{
if let rustfs_common::heal_channel::HealChannelCommand::Start { response_tx, .. } = command {
heal_starts_clone.fetch_add(1, Ordering::Relaxed);
healed_versions_clone
.lock()
.expect("heal version capture lock should not be poisoned")
.push(request.object_version_id);
let _ = response_tx.send(Ok(HealAdmissionResult::Accepted));
}
}
@@ -1813,18 +1739,13 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
let bucket = "src-archive";
let object = "snapshots/37b3f20d941e2f5e6d99114d9bb2f3e67a8a2e5c9c4c5a1b0d6e7f8091a2b3c4";
let orphan_version = Uuid::from_u128(0x1934);
let shared_version = Uuid::from_u128(0x1935);
let orphan_metadata = metadata_for_object_version(bucket, object, Some(orphan_version));
let shared_metadata = metadata_for_object_version(bucket, object, Some(shared_version));
write_test_object_metadata_bytes(&temp_dir, bucket, object, &orphan_metadata).await;
let mut expected_metadata = vec![(temp_dir.join(bucket).join(object).join("xl.meta"), orphan_metadata.clone())];
let metadata = metadata_for_object(bucket, object);
write_test_object_metadata_bytes(&temp_dir, bucket, object, &metadata).await;
let mut disks = vec![scanner.local_disk.clone()];
for disk_name in ["disk2", "disk3", "disk4"] {
let disk_root = temp_dir.join(disk_name);
write_test_object_metadata_bytes(&disk_root, bucket, object, &shared_metadata).await;
expected_metadata.push((disk_root.join(bucket).join(object).join("xl.meta"), shared_metadata.clone()));
write_test_object_metadata_bytes(&disk_root, bucket, object, &metadata).await;
let endpoint = Endpoint::try_from(disk_root.to_string_lossy().as_ref()).expect("failed to create extra disk endpoint");
let disk = new_disk(
&endpoint,
@@ -1873,29 +1794,8 @@ async fn test_scan_folder_exits_when_abandoned_child_listing_finishes() {
.new_cache
.checked_flatten(bucket)
.expect("healed cache must contain canonical child links");
// The fixture intentionally exposes two divergent version histories, so
// the scanner keeps both logical versions visible while discovering heals.
assert_eq!(root.objects, 2);
assert_eq!(root.objects, 1);
assert!(heal_starts.load(Ordering::Relaxed) > 0, "test must execute the heal child-link path");
let orphan_version_text = orphan_version.to_string();
assert!(
healed_versions
.lock()
.expect("heal version capture lock should not be poisoned")
.iter()
.any(|version| version.as_deref() == Some(orphan_version_text.as_str())),
"sub-quorum orphan version must be submitted as an exact heal candidate"
);
for (path, expected) in expected_metadata {
assert_eq!(
tokio::fs::read(&path)
.await
.expect("scanner discovery must not delete metadata"),
expected,
"scanner discovery must not modify candidate metadata: {}",
path.display()
);
}
}
#[tokio::test]
+11 -4
View File
@@ -158,10 +158,11 @@ added by backlog#1153 infra-4.
## Coverage
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).
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).
- **CI**: `.github/workflows/coverage.yml` runs every Sunday and on manual
dispatch: `cargo llvm-cov nextest --workspace --exclude e2e_test` under the
@@ -174,6 +175,12 @@ baseline later (backlog#1153 infra-6, report-only first).
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'
+264
View File
@@ -0,0 +1,264 @@
#!/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())
+27 -14
View File
@@ -47,6 +47,31 @@ 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)
@@ -54,24 +79,12 @@ 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:
data = export["data"][0]
files = data["files"]
totals = data["totals"]["lines"]
except (KeyError, IndexError) as exc:
crates, totals = load_coverage(path, root)
except (KeyError, IndexError, ValueError) 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,
+1 -1
View File
@@ -241,7 +241,7 @@ env \
RUSTFS_TEST_VAULT_FAILOVER_MARKER="$MARKER" \
RUSTFS_TEST_VAULT_OLD_LEADER="$OLD_LEADER" \
cargo test -p rustfs-kms --test vault_ha_failover_live \
vault_raft_leader_failure_preserves_kv2_and_transit_decrypts -- \
vault_raft_leader_failure_recovers_kv2_and_transit_decrypts -- \
--ignored --nocapture --test-threads=1 &
TEST_PID=$!