fix(ci): restore main checks (#6469)

* fix(ci): restore main checks

* fix(scanner): bootstrap pristine usage state

* fix(scanner): reject empty usage snapshots
This commit is contained in:
Zhengchao An
2026-08-24 09:33:00 +08:00
committed by GitHub
parent ebff02304d
commit 99fb77b164
3 changed files with 245 additions and 78 deletions
+32 -9
View File
@@ -1318,8 +1318,8 @@ pub(super) async fn persisted_usage_floor(
}; };
for primary_path in [DATA_USAGE_OBJ_NAME_PATH.as_str(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()] { for primary_path in [DATA_USAGE_OBJ_NAME_PATH.as_str(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()] {
let backup_path = format!("{primary_path}.bkp"); let backup_path = format!("{primary_path}.bkp");
let primary_epoch = match read_config(storeapi.clone(), primary_path).await { let primary_epoch = match read_config_with_revision(storeapi.clone(), primary_path).await {
Ok(data) => { Ok((Some(data), _)) => {
let usage = serde_json::from_slice::<DataUsageInfo>(&data).map_err(|err| { let usage = serde_json::from_slice::<DataUsageInfo>(&data).map_err(|err| {
ScannerError::Other(format!("failed to decode scanner usage floor from {primary_path}: {err}")) ScannerError::Other(format!("failed to decode scanner usage floor from {primary_path}: {err}"))
})?; })?;
@@ -1332,7 +1332,7 @@ pub(super) async fn persisted_usage_floor(
update_floor(&mut floor, &usage, primary_path)?; update_floor(&mut floor, &usage, primary_path)?;
Some(epoch) Some(epoch)
} }
Err(EcstoreError::ConfigNotFound) => None, Ok((None, _)) => None,
Err(err) => { Err(err) => {
return Err(ScannerError::Other(format!( return Err(ScannerError::Other(format!(
"failed to read scanner usage epoch floor from {primary_path}: {err}" "failed to read scanner usage epoch floor from {primary_path}: {err}"
@@ -1340,8 +1340,8 @@ pub(super) async fn persisted_usage_floor(
} }
}; };
let mut any_found = primary_epoch.is_some(); let mut any_found = primary_epoch.is_some();
match read_config(storeapi.clone(), &backup_path).await { match read_config_with_revision(storeapi.clone(), &backup_path).await {
Ok(data) => { Ok((Some(data), _)) => {
any_found = true; any_found = true;
let usage = serde_json::from_slice::<DataUsageInfo>(&data).map_err(|err| { let usage = serde_json::from_slice::<DataUsageInfo>(&data).map_err(|err| {
ScannerError::Other(format!("failed to decode scanner usage floor from {backup_path}: {err}")) ScannerError::Other(format!("failed to decode scanner usage floor from {backup_path}: {err}"))
@@ -1359,7 +1359,7 @@ pub(super) async fn persisted_usage_floor(
update_floor(&mut floor, &usage, &backup_path)?; update_floor(&mut floor, &usage, &backup_path)?;
} }
} }
Err(EcstoreError::ConfigNotFound) => {} Ok((None, _)) => {}
Err(err) => { Err(err) => {
return Err(ScannerError::Other(format!( return Err(ScannerError::Other(format!(
"failed to read scanner usage epoch floor from {backup_path}: {err}" "failed to read scanner usage epoch floor from {backup_path}: {err}"
@@ -1373,9 +1373,32 @@ pub(super) async fn persisted_usage_floor(
} }
if !found_any { if !found_any {
return Err(ScannerError::Other( let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else {
"persisted scanner usage floor has no authoritative baseline".to_string(), return Err(ScannerError::Other(
)); "scanner usage floor changed before pristine state confirmation".to_string(),
));
};
for path in [
DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
] {
match read_config_with_revision(storeapi.clone(), &path).await {
Ok((None, _)) => {}
Ok((Some(_), _)) => {
return Err(ScannerError::Other(format!(
"scanner usage floor changed while confirming pristine state: {path} appeared"
)));
}
Err(err) => {
return Err(ScannerError::Other(format!(
"failed to confirm pristine scanner usage floor at {path}: {err}"
)));
}
}
}
drop(publication_admission);
} }
let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi, read_epoch).await else { let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi, read_epoch).await else {
return Err(ScannerError::Other( return Err(ScannerError::Other(
+64 -47
View File
@@ -104,6 +104,38 @@ pub(super) async fn usage_snapshot_for_epoch_fence(
Ok(None) Ok(None)
} }
async fn read_usage_snapshot_for_epoch_fence(
storeapi: Arc<impl ScannerObjectIO>,
) -> Result<(Option<DataUsageInfo>, DataUsageCacheRevision), ScannerError> {
let (primary, revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to read scanner usage epoch fence: {err}")))?;
let usage = usage_snapshot_for_epoch_fence(storeapi, primary.as_deref()).await?;
Ok((usage, revision))
}
async fn confirm_usage_snapshot_absent_for_bootstrap(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
expected_publication_epoch: u64,
) -> Result<crate::ScannerDataUsagePublicationAdmission, ScannerError> {
let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_publication_epoch).await
else {
return Err(ScannerError::Other(
"scanner publication epoch changed before confirming pristine usage state".to_string(),
));
};
let (usage, _) = read_usage_snapshot_for_epoch_fence(storeapi.clone()).await?;
if usage.is_some() {
return Err(ScannerError::Other(
"scanner usage baseline appeared while confirming pristine bootstrap".to_string(),
));
}
drop(publication_admission);
scanner_publication_admission_for_epoch(storeapi, expected_publication_epoch)
.await
.ok_or_else(|| ScannerError::Other("scanner publication epoch changed during pristine usage confirmation".to_string()))
}
pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
ctx: &CancellationToken, ctx: &CancellationToken,
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>, storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
@@ -128,19 +160,10 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
"scanner usage epoch fence changed while recovery reset was in progress".to_string(), "scanner usage epoch fence changed while recovery reset was in progress".to_string(),
)); ));
} }
let (primary, revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) let (usage, revision) = read_usage_snapshot_for_epoch_fence(storeapi.clone()).await?;
.await let Some(mut usage) = usage else {
.map_err(|err| ScannerError::Other(format!("failed to read scanner usage epoch fence: {err}")))?; let _publication_admission = confirm_usage_snapshot_absent_for_bootstrap(storeapi, read_epoch).await?;
let Some(mut usage) = usage_snapshot_for_epoch_fence(storeapi.clone(), primary.as_deref()).await? else { return Ok(());
let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else {
if retry < SCANNER_PERSIST_CAS_RETRIES {
continue;
}
return Err(ScannerError::Other(
"scanner usage epoch fence changed while confirming a missing usage baseline".to_string(),
));
};
return Err(ScannerError::Other("authoritative scanner usage baseline is missing".to_string()));
}; };
match usage.scanner_epoch { match usage.scanner_epoch {
Some(epoch) if epoch > claimed_epoch => { Some(epoch) if epoch > claimed_epoch => {
@@ -274,8 +297,9 @@ pub(super) async fn claim_scanner_leadership(
let Some(read_epoch) = scanner_publication_epoch(storeapi.clone()).await else { let Some(read_epoch) = scanner_publication_epoch(storeapi.clone()).await else {
return false; return false;
}; };
let (usage_primary, _) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await { let usage_baseline_missing = match read_usage_snapshot_for_epoch_fence(storeapi.clone()).await {
Ok(result) => result, Ok((Some(_), _)) => false,
Ok((None, _)) => true,
Err(err) => { Err(err) => {
error!( error!(
target: "rustfs::scanner", target: "rustfs::scanner",
@@ -290,40 +314,33 @@ pub(super) async fn claim_scanner_leadership(
return false; return false;
} }
}; };
match usage_snapshot_for_epoch_fence(storeapi.clone(), usage_primary.as_deref()).await {
Ok(Some(_)) => {}
Ok(None) => {
warn!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %DATA_USAGE_OBJ_NAME_PATH.as_str(),
state = "leader_usage_baseline_missing",
"Scanner leadership claim deferred until a usage baseline is published"
);
return false;
}
Err(err) => {
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
path = %DATA_USAGE_OBJ_NAME_PATH.as_str(),
state = "leader_usage_baseline_invalid",
error = %err,
"Scanner leadership claim deferred because the usage baseline is invalid"
);
return false;
}
}
let save_result = { let save_result = {
let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { let _publication_admission = if usage_baseline_missing {
if retry < SCANNER_PERSIST_CAS_RETRIES { match confirm_usage_snapshot_absent_for_bootstrap(storeapi.clone(), read_epoch).await {
continue; Ok(publication_admission) => publication_admission,
Err(err) => {
error!(
target: "rustfs::scanner",
event = EVENT_SCANNER_PERSIST_STATE,
component = LOG_COMPONENT_SCANNER,
subsystem = LOG_SUBSYSTEM_RUNTIME,
state = "leader_usage_baseline_changed",
path = %DATA_USAGE_OBJ_NAME_PATH.as_str(),
error = %err,
"Scanner leadership claim deferred because the pristine usage state changed"
);
return false;
}
} }
return false; } else {
let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await
else {
if retry < SCANNER_PERSIST_CAS_RETRIES {
continue;
}
return false;
};
publication_admission
}; };
save_config_with_preconditions(storeapi.clone(), &DATA_USAGE_BLOOM_NAME_PATH, data.clone(), revision.preconditions()) save_config_with_preconditions(storeapi.clone(), &DATA_USAGE_BLOOM_NAME_PATH, data.clone(), revision.preconditions())
.await .await
+149 -22
View File
@@ -23,7 +23,7 @@ use crate::{
}; };
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet};
use std::io::Cursor; use std::io::Cursor;
use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::task::Poll; use std::task::Poll;
use temp_env::{with_var, with_var_unset}; use temp_env::{with_var, with_var_unset};
use tokio::io::AsyncReadExt; use tokio::io::AsyncReadExt;
@@ -336,6 +336,7 @@ impl Drop for ScannerDefaultCycleGuard {
struct MemoryConfigStore { struct MemoryConfigStore {
objects: Mutex<HashMap<String, Vec<u8>>>, objects: Mutex<HashMap<String, Vec<u8>>>,
revisions: Mutex<HashMap<String, u64>>, revisions: Mutex<HashMap<String, u64>>,
insert_after_gets: Mutex<HashMap<String, Vec<u8>>>,
non_regular_objects: Mutex<HashSet<String>>, non_regular_objects: Mutex<HashSet<String>>,
fail_put_number: Mutex<HashMap<String, usize>>, fail_put_number: Mutex<HashMap<String, usize>>,
object_not_found_put_number: Mutex<HashMap<String, usize>>, object_not_found_put_number: Mutex<HashMap<String, usize>>,
@@ -346,12 +347,36 @@ struct MemoryConfigStore {
replace_after_successful_puts: Mutex<HashMap<String, (usize, Vec<u8>)>>, replace_after_successful_puts: Mutex<HashMap<String, (usize, Vec<u8>)>>,
put_counts: Mutex<HashMap<String, usize>>, put_counts: Mutex<HashMap<String, usize>>,
publication_admission_blocked: AtomicBool, publication_admission_blocked: AtomicBool,
block_publication_after_admissions: AtomicUsize,
} }
fn memory_config_key(bucket: &str, object: &str) -> String { fn memory_config_key(bucket: &str, object: &str) -> String {
format!("{bucket}/{object}") format!("{bucket}/{object}")
} }
async fn insert_usage_after_first_legacy_backup_read(store: &MemoryConfigStore) {
let legacy_backup = format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str());
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
usage.scanner_epoch = Some(7);
usage.scanner_cycle = Some(11);
store.insert_after_gets.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, &legacy_backup),
serde_json::to_vec(&usage).expect("usage snapshot should encode"),
);
}
async fn claim_test_scanner_leadership(store: Arc<MemoryConfigStore>) -> (bool, u64) {
let ctx = CancellationToken::new();
let mut revision = DataUsageCacheRevision::Missing;
let mut cycle = CurrentCycle {
next: 12,
..Default::default()
};
let mut persisted_epoch = 0;
let claimed = claim_scanner_leadership(&ctx, store, &mut cycle, &mut revision, &mut persisted_epoch).await;
(claimed, persisted_epoch)
}
#[async_trait::async_trait] #[async_trait::async_trait]
impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore { impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore {
type Error = EcstoreError; type Error = EcstoreError;
@@ -371,13 +396,21 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore {
_opts: &ObjectOptions, _opts: &ObjectOptions,
) -> EcstoreResult<GetObjectReader> { ) -> EcstoreResult<GetObjectReader> {
let key = memory_config_key(bucket, object); let key = memory_config_key(bucket, object);
let data = self let inserted_data = self.insert_after_gets.lock().await.remove(&key);
.objects let data = {
.lock() let mut objects = self.objects.lock().await;
.await let data = objects.get(&key).cloned();
.get(&key) if let Some(inserted_data) = inserted_data.as_ref() {
.cloned() objects.insert(key.clone(), inserted_data.clone());
.ok_or(EcstoreError::FileNotFound)?; }
data
};
if inserted_data.is_some() {
let mut revisions = self.revisions.lock().await;
let revision = revisions.get(&key).copied().unwrap_or(0) + 1;
revisions.insert(key.clone(), revision);
}
let data = data.ok_or(EcstoreError::FileNotFound)?;
let data_len = i64::try_from(data.len()).expect("memory test object length should fit in i64"); let data_len = i64::try_from(data.len()).expect("memory test object length should fit in i64");
let revision = *self.revisions.lock().await.entry(key.clone()).or_insert(1); let revision = *self.revisions.lock().await.entry(key.clone()).or_insert(1);
let is_dir = self.non_regular_objects.lock().await.contains(&key); let is_dir = self.non_regular_objects.lock().await.contains(&key);
@@ -2033,8 +2066,6 @@ async fn scanner_startup_prefers_v2_over_legacy_usage() {
#[tokio::test] #[tokio::test]
async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state() { async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());
assert!(persisted_usage_floor(store.clone()).await.is_err());
store.objects.lock().await.insert( store.objects.lock().await.insert(
memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()),
b"not-json".to_vec(), b"not-json".to_vec(),
@@ -2062,6 +2093,67 @@ async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state()
assert!(persisted_usage_floor(store).await.is_err()); assert!(persisted_usage_floor(store).await.is_err());
} }
#[tokio::test]
async fn scanner_usage_floor_fails_closed_on_zero_byte_usage_objects() {
for path in [
DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
] {
let key = memory_config_key(RUSTFS_META_BUCKET, &path);
let existing = Arc::new(MemoryConfigStore::default());
existing.objects.lock().await.insert(key.clone(), Vec::new());
let err = persisted_usage_floor(existing)
.await
.expect_err("an empty usage object must not be treated as missing");
assert!(
err.to_string()
.contains(&format!("failed to decode scanner usage floor from {path}:")),
"unexpected error for {path}: {err}"
);
let appearing = Arc::new(MemoryConfigStore::default());
appearing.insert_after_gets.lock().await.insert(key, Vec::new());
let err = persisted_usage_floor(appearing)
.await
.expect_err("an empty usage object appearing during confirmation must prevent pristine bootstrap");
assert!(
err.to_string().contains("changed while confirming pristine state"),
"unexpected confirmation error for {path}: {err}"
);
}
}
#[tokio::test]
async fn scanner_usage_floor_requires_publication_admission_for_pristine_bootstrap() {
let store = Arc::new(MemoryConfigStore::default());
store.publication_admission_blocked.store(true, Ordering::Release);
assert!(persisted_usage_floor(store).await.is_err());
}
#[tokio::test]
async fn scanner_usage_floor_fails_closed_when_usage_appears_during_pristine_confirmation() {
let store = Arc::new(MemoryConfigStore::default());
insert_usage_after_first_legacy_backup_read(store.as_ref()).await;
let err = persisted_usage_floor(store)
.await
.expect_err("an appearing usage snapshot must prevent pristine bootstrap");
assert!(err.to_string().contains("changed while confirming pristine state"));
}
#[tokio::test]
async fn scanner_usage_floor_rejects_publication_change_during_pristine_confirmation() {
let store = Arc::new(MemoryConfigStore::default());
store.block_publication_after_admissions.store(2, Ordering::Release);
assert!(persisted_usage_floor(store).await.is_err());
}
#[tokio::test] #[tokio::test]
async fn scanner_usage_backup_uses_durable_cycle_cadence_across_tasks() { async fn scanner_usage_backup_uses_durable_cycle_cadence_across_tasks() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());
@@ -2156,7 +2248,17 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
} }
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> { async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
(!self.publication_admission_blocked.load(Ordering::Acquire)).then(crate::ScannerDataUsagePublicationAdmission::unfenced) if self.publication_admission_blocked.load(Ordering::Acquire) {
return None;
}
if self
.block_publication_after_admissions
.fetch_update(Ordering::AcqRel, Ordering::Acquire, |remaining| remaining.checked_sub(1))
== Ok(1)
{
self.publication_admission_blocked.store(true, Ordering::Release);
}
Some(crate::ScannerDataUsagePublicationAdmission::unfenced())
} }
} }
@@ -2363,21 +2465,46 @@ async fn test_leadership_claim_rejects_terminal_epoch() {
} }
#[tokio::test] #[tokio::test]
async fn leadership_claim_defers_without_usage_baseline_before_bloom_write() { async fn scanner_bootstraps_leadership_when_usage_snapshots_are_stably_absent() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());
let ctx = CancellationToken::new(); let floor = persisted_usage_floor(store.clone())
let mut revision = DataUsageCacheRevision::Missing; .await
let mut cycle = CurrentCycle { .expect("pristine usage state should provide the initial floor");
next: 12, assert_eq!(floor, PersistedUsageFloor::default());
..Default::default() let (claimed, persisted_epoch) = claim_test_scanner_leadership(store.clone()).await;
}; assert!(claimed);
let mut persisted_epoch = 0; let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH)
.await
assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch,).await); .expect("pristine leadership claim should persist");
assert!(read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); let (persisted_cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("persisted leadership claim should decode");
assert_eq!(persisted_cycle.next, 12);
assert_eq!(leader_epoch, 1);
assert_eq!(persisted_epoch, 1);
assert!(read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()).await.is_err()); assert!(read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()).await.is_err());
} }
#[tokio::test]
async fn leadership_claim_fails_closed_when_usage_appears_during_pristine_confirmation() {
let store = Arc::new(MemoryConfigStore::default());
insert_usage_after_first_legacy_backup_read(store.as_ref()).await;
let (claimed, persisted_epoch) = claim_test_scanner_leadership(store.clone()).await;
assert!(!claimed);
assert_eq!(persisted_epoch, 0);
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err());
}
#[tokio::test]
async fn leadership_claim_rejects_publication_change_during_pristine_confirmation() {
let store = Arc::new(MemoryConfigStore::default());
store.block_publication_after_admissions.store(2, Ordering::Release);
let (claimed, persisted_epoch) = claim_test_scanner_leadership(store.clone()).await;
assert!(!claimed);
assert_eq!(persisted_epoch, 0);
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err());
}
#[tokio::test] #[tokio::test]
async fn leadership_claim_defers_on_corrupt_usage_baseline_without_bloom_write() { async fn leadership_claim_defers_on_corrupt_usage_baseline_without_bloom_write() {
let store = Arc::new(MemoryConfigStore::default()); let store = Arc::new(MemoryConfigStore::default());