diff --git a/crates/scanner/src/scanner/cycle_state.rs b/crates/scanner/src/scanner/cycle_state.rs index db2e46ff8..03a6b03ec 100644 --- a/crates/scanner/src/scanner/cycle_state.rs +++ b/crates/scanner/src/scanner/cycle_state.rs @@ -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()] { let backup_path = format!("{primary_path}.bkp"); - let primary_epoch = match read_config(storeapi.clone(), primary_path).await { - Ok(data) => { + let primary_epoch = match read_config_with_revision(storeapi.clone(), primary_path).await { + Ok((Some(data), _)) => { let usage = serde_json::from_slice::(&data).map_err(|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)?; Some(epoch) } - Err(EcstoreError::ConfigNotFound) => None, + Ok((None, _)) => None, Err(err) => { return Err(ScannerError::Other(format!( "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(); - match read_config(storeapi.clone(), &backup_path).await { - Ok(data) => { + match read_config_with_revision(storeapi.clone(), &backup_path).await { + Ok((Some(data), _)) => { any_found = true; let usage = serde_json::from_slice::(&data).map_err(|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)?; } } - Err(EcstoreError::ConfigNotFound) => {} + Ok((None, _)) => {} Err(err) => { return Err(ScannerError::Other(format!( "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 { - return Err(ScannerError::Other( - "persisted scanner usage floor has no authoritative baseline".to_string(), - )); + let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { + 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 { return Err(ScannerError::Other( diff --git a/crates/scanner/src/scanner/leadership.rs b/crates/scanner/src/scanner/leadership.rs index 8183b7602..69d54f2aa 100644 --- a/crates/scanner/src/scanner/leadership.rs +++ b/crates/scanner/src/scanner/leadership.rs @@ -104,6 +104,38 @@ pub(super) async fn usage_snapshot_for_epoch_fence( Ok(None) } +async fn read_usage_snapshot_for_epoch_fence( + storeapi: Arc, +) -> Result<(Option, 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, + expected_publication_epoch: u64, +) -> Result { + 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( ctx: &CancellationToken, storeapi: Arc, @@ -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(), )); } - 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 Some(mut usage) = usage_snapshot_for_epoch_fence(storeapi.clone(), primary.as_deref()).await? else { - 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())); + let (usage, revision) = read_usage_snapshot_for_epoch_fence(storeapi.clone()).await?; + let Some(mut usage) = usage else { + let _publication_admission = confirm_usage_snapshot_absent_for_bootstrap(storeapi, read_epoch).await?; + return Ok(()); }; match usage.scanner_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 { return false; }; - let (usage_primary, _) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await { - Ok(result) => result, + let usage_baseline_missing = match read_usage_snapshot_for_epoch_fence(storeapi.clone()).await { + Ok((Some(_), _)) => false, + Ok((None, _)) => true, Err(err) => { error!( target: "rustfs::scanner", @@ -290,40 +314,33 @@ pub(super) async fn claim_scanner_leadership( 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 Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { - if retry < SCANNER_PERSIST_CAS_RETRIES { - continue; + let _publication_admission = if usage_baseline_missing { + match confirm_usage_snapshot_absent_for_bootstrap(storeapi.clone(), read_epoch).await { + 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()) .await diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index b59cdc2d8..983c93b1b 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -23,7 +23,7 @@ use crate::{ }; use std::collections::{HashMap, HashSet}; use std::io::Cursor; -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::task::Poll; use temp_env::{with_var, with_var_unset}; use tokio::io::AsyncReadExt; @@ -336,6 +336,7 @@ impl Drop for ScannerDefaultCycleGuard { struct MemoryConfigStore { objects: Mutex>>, revisions: Mutex>, + insert_after_gets: Mutex>>, non_regular_objects: Mutex>, fail_put_number: Mutex>, object_not_found_put_number: Mutex>, @@ -346,12 +347,36 @@ struct MemoryConfigStore { replace_after_successful_puts: Mutex)>>, put_counts: Mutex>, publication_admission_blocked: AtomicBool, + block_publication_after_admissions: AtomicUsize, } fn memory_config_key(bucket: &str, object: &str) -> String { 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) -> (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] impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore { type Error = EcstoreError; @@ -371,13 +396,21 @@ impl crate::storage_api::scanner_io::ObjectIO for MemoryConfigStore { _opts: &ObjectOptions, ) -> EcstoreResult { let key = memory_config_key(bucket, object); - let data = self - .objects - .lock() - .await - .get(&key) - .cloned() - .ok_or(EcstoreError::FileNotFound)?; + let inserted_data = self.insert_after_gets.lock().await.remove(&key); + let data = { + let mut objects = self.objects.lock().await; + let data = objects.get(&key).cloned(); + if let Some(inserted_data) = inserted_data.as_ref() { + objects.insert(key.clone(), inserted_data.clone()); + } + 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 revision = *self.revisions.lock().await.entry(key.clone()).or_insert(1); 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] async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state() { let store = Arc::new(MemoryConfigStore::default()); - assert!(persisted_usage_floor(store.clone()).await.is_err()); - store.objects.lock().await.insert( memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), 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()); } +#[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] async fn scanner_usage_backup_uses_durable_cycle_cadence_across_tasks() { let store = Arc::new(MemoryConfigStore::default()); @@ -2156,7 +2248,17 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore { } async fn scanner_data_usage_publication_admission(&self) -> Option { - (!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] -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 ctx = CancellationToken::new(); - let mut revision = DataUsageCacheRevision::Missing; - let mut cycle = CurrentCycle { - next: 12, - ..Default::default() - }; - let mut persisted_epoch = 0; - - assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch,).await); - assert!(read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); + let floor = persisted_usage_floor(store.clone()) + .await + .expect("pristine usage state should provide the initial floor"); + assert_eq!(floor, PersistedUsageFloor::default()); + let (claimed, persisted_epoch) = claim_test_scanner_leadership(store.clone()).await; + assert!(claimed); + let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) + .await + .expect("pristine leadership claim should persist"); + 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()); } +#[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] async fn leadership_claim_defers_on_corrupt_usage_baseline_without_bloom_write() { let store = Arc::new(MemoryConfigStore::default());