From 170a4c76407e83aee12fc4b1bf4dc496e0ac82c5 Mon Sep 17 00:00:00 2001 From: Zhengchao An Date: Mon, 24 Aug 2026 14:29:09 +0800 Subject: [PATCH] fix(scanner): bootstrap pristine usage baseline (#6471) * fix(ecstore): fence pool metadata replica updates * fix(ecstore): block decommission on unsafe pool metadata * fix(ecstore): block writes after pool metadata save errors * fix(ecstore): latch pool metadata writes before await * fix(scanner): bootstrap pristine usage baseline --- crates/data-usage/src/data_usage.rs | 10 +- crates/ecstore/src/core/pools.rs | 6 +- crates/ecstore/src/data_usage/mod.rs | 33 ++- crates/scanner/src/scanner.rs | 141 ++++++++++-- crates/scanner/src/scanner/cycle_state.rs | 70 +++++- crates/scanner/src/scanner/leadership.rs | 213 ++++++++++++------ crates/scanner/src/scanner/tests.rs | 253 ++++++++++++++++++---- 7 files changed, 582 insertions(+), 144 deletions(-) diff --git a/crates/data-usage/src/data_usage.rs b/crates/data-usage/src/data_usage.rs index f26410329..68a4d9822 100644 --- a/crates/data-usage/src/data_usage.rs +++ b/crates/data-usage/src/data_usage.rs @@ -573,6 +573,10 @@ pub struct DataUsageInfo { /// cycle (or retained a compatible last-known-good cache). #[serde(default)] pub usage_snapshot_partial: bool, + /// Durable marker for a first-start bootstrap that has not produced an + /// authoritative usage snapshot yet. + #[serde(default, skip_serializing_if = "std::ops::Not::not")] + pub usage_snapshot_bootstrap_pending: bool, /// Deprecated kept here for backward compatibility reasons pub bucket_sizes: HashMap, /// Per-disk snapshot information when available @@ -1962,7 +1966,8 @@ impl DataUsageInfo { /// Whether this snapshot authoritatively covers every reported bucket. pub fn is_complete_bucket_usage_snapshot(&self) -> bool { - self.usage_snapshot_complete + !self.usage_snapshot_bootstrap_pending + && self.usage_snapshot_complete && self.last_update.is_some() && u64::try_from(self.buckets_usage.len()).ok() == Some(self.buckets_count) } @@ -1971,7 +1976,8 @@ impl DataUsageInfo { /// admin display. Partial data is accepted only with unique set states, /// a plan digest for every state, and at least one usable generation. pub fn is_valid_partial_snapshot(&self) -> bool { - if !self.usage_snapshot_partial + if self.usage_snapshot_bootstrap_pending + || !self.usage_snapshot_partial || self.usage_snapshot_converged != Some(false) || self.last_update.is_none() || self.scanner_cycle.is_none() diff --git a/crates/ecstore/src/core/pools.rs b/crates/ecstore/src/core/pools.rs index 1942c9514..26bb1066e 100644 --- a/crates/ecstore/src/core/pools.rs +++ b/crates/ecstore/src/core/pools.rs @@ -33,7 +33,7 @@ use crate::bucket::{ use crate::cache_value::metacache_set::{ListPathRawOptions, list_path_raw}; use crate::config::com::{ CONFIG_PREFIX, delete_config, read_config_limited_preserve_empty, read_config_limited_preserve_empty_with_metadata, - read_config_no_lock_preserve_empty_with_metadata, read_config_preserve_empty, save_config, save_config_with_opts, + read_config_no_lock_preserve_empty_with_metadata, read_config_preserve_empty, save_config_with_opts, save_config_with_opts_quiet, }; use crate::data_movement; @@ -4848,7 +4848,7 @@ impl ECStore { let mut pool_meta = self.pool_meta.write().await; record_decommission_unresolved_entry(&mut pool_meta, idx, generation, entry)?; } - self.save_current_pool_meta() + self.save_current_pool_meta(&[idx]) .await .map_err(|err| Error::other(format!("decommission unresolved entry ledger save failed: {err}"))) } @@ -9038,7 +9038,7 @@ impl ECStore { self.run_guarded_decommission_side_effect(rx, &operation_gate, || { self.check_after_decommission_unfenced(idx, generation) }) - .await + .await } async fn check_after_decommission_unfenced( diff --git a/crates/ecstore/src/data_usage/mod.rs b/crates/ecstore/src/data_usage/mod.rs index 37f14c4ed..c1db05676 100644 --- a/crates/ecstore/src/data_usage/mod.rs +++ b/crates/ecstore/src/data_usage/mod.rs @@ -743,7 +743,12 @@ where let backup_seed = load_data_usage_for_bucket_removal(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) .await? .map(|(data_usage_info, _)| data_usage_info) - .or_else(|| primary_seed.clone()); + .filter(|data_usage_info| !data_usage_info.usage_snapshot_bootstrap_pending) + .or_else(|| { + primary_seed + .clone() + .filter(|data_usage_info| !data_usage_info.usage_snapshot_bootstrap_pending) + }); remove_bucket_usage_from_object_with_retries_and_publication( store, DATA_USAGE_OBJ_BACKUP_PATH.as_str(), @@ -4646,6 +4651,32 @@ mod tests { assert!(state.backup_object.is_none()); } + #[tokio::test] + async fn remove_bucket_usage_does_not_seed_backup_from_pristine_bootstrap_marker() { + let marker = DataUsageInfo { + last_update: Some(SystemTime::now()), + usage_snapshot_converged: Some(false), + usage_snapshot_bootstrap_pending: true, + ..Default::default() + }; + let store = Arc::new(UsageCasStore { + state: Mutex::new(UsageCasState { + object: Some((serde_json::to_vec(&marker).expect("bootstrap marker should encode"), 1)), + ..Default::default() + }), + }); + + remove_bucket_usage_from_backend_with_store(store.as_ref(), "bucket-a") + .await + .expect("bucket removal should preserve the pending primary without creating a backup"); + + let state = store.state.lock().await; + assert!(state.backup_object.is_none()); + let saved = serde_json::from_slice::(&state.object.as_ref().expect("pending primary should remain").0) + .expect("pending primary should decode"); + assert!(saved.usage_snapshot_bootstrap_pending); + } + #[tokio::test] async fn remove_bucket_usage_migrates_legacy_snapshot_without_hiding_other_buckets() { let mut legacy = data_usage_info_for_test("bucket-a", 2, 84, SystemTime::now()); diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index f7738c838..bc9cc8296 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -366,7 +366,8 @@ pub(super) fn data_usage_info_has_persisted_baseline_identity(info: &DataUsageIn // complete: a timestamp, a scanner cycle, and an exact bucket cardinality. // A current snapshot with only scanner_epoch/scanner_cycle (or an explicit // incomplete marker) is not evidence of a durable usage baseline. - !info.usage_snapshot_complete + !info.usage_snapshot_bootstrap_pending + && !info.usage_snapshot_complete && info.scanner_epoch.is_none() && info.usage_snapshot_converged != Some(false) && info.last_update.is_some() @@ -374,6 +375,21 @@ pub(super) fn data_usage_info_has_persisted_baseline_identity(info: &DataUsageIn && u64::try_from(info.buckets_usage.len()).ok() == Some(info.buckets_count) } +pub(super) fn data_usage_info_is_pristine_bootstrap_pending(info: &DataUsageInfo) -> bool { + if info.last_update.is_none() || info.scanner_cycle.is_some() { + return false; + } + + let expected = DataUsageInfo { + last_update: info.last_update, + scanner_epoch: info.scanner_epoch, + usage_snapshot_converged: Some(false), + usage_snapshot_bootstrap_pending: true, + ..Default::default() + }; + info == &expected +} + fn usage_cache_needs_prompt_scan(authoritative: &DataUsageInfo, observed: Option<&DataUsageInfo>) -> bool { data_usage_info_is_cold(authoritative) || observed.is_some_and(|observed| observed_data_usage_is_newer(observed, authoritative)) @@ -637,6 +653,29 @@ async fn initial_scanner_startup_usage_state(storeapi: &Arc) -> (bool, (persisted_usage_cache_is_cold_for_startup(storeapi).await, has_buckets) } +fn scanner_cycle_state_is_pristine( + cycle_info: &CurrentCycle, + leader_epoch: u64, + cycle_revision: &DataUsageCacheRevision, +) -> bool { + cycle_info.next == 0 && leader_epoch == 0 && matches!(cycle_revision, DataUsageCacheRevision::Missing) +} + +fn scanner_may_bootstrap_missing_usage_floor( + cycle_info: &CurrentCycle, + leader_epoch: u64, + cycle_revision: &DataUsageCacheRevision, +) -> bool { + // The server becomes ready before the scanner starts, so a first bucket may + // already exist. The bootstrap marker is non-authoritative; only prior + // durable scanner progress must block its creation. + scanner_cycle_state_is_pristine(cycle_info, leader_epoch, cycle_revision) +} + +fn scanner_may_resume_pristine_usage_bootstrap(cycle_info: &CurrentCycle) -> bool { + cycle_info.next == 0 +} + pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc) { let (startup_features, startup_maintenance_generation) = configure_scanner_defaults(&ctx, &storeapi).await; // Force init global sleeper so config is read once at startup. @@ -1153,7 +1192,7 @@ where LockLost: Future, { let fence_ctx = ctx.child_token(); - let claim = claim_scanner_leadership(&fence_ctx, storeapi, cycle_info, cycle_revision, leader_epoch); + let claim = claim_scanner_leadership(&fence_ctx, storeapi, cycle_info, cycle_revision, leader_epoch, false); tokio::pin!(claim); tokio::pin!(lock_lost); tokio::select! { @@ -2100,24 +2139,81 @@ async fn run_data_scanner_with_maintenance_state( return Err(err); } }; - let usage_floor = match persisted_usage_floor(storeapi.clone()).await { - Ok(floor) => floor, - 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 = "usage_floor_load_failed", - error = %err, - "Scanner stopped because the persisted usage floor could not be loaded" - ); - global_metrics().set_cycle(None).await; - return Ok(()); + let may_bootstrap_missing_usage_floor = scanner_may_bootstrap_missing_usage_floor(&cycle_info, leader_epoch, &cycle_revision); + let (usage_floor, usage_floor_startup) = + match persisted_usage_floor_for_startup(storeapi.clone(), may_bootstrap_missing_usage_floor).await { + Ok(result) => result, + 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 = "usage_floor_load_failed", + error = %err, + "Scanner stopped because the persisted usage floor could not be loaded" + ); + global_metrics().set_cycle(None).await; + return Ok(()); + } + }; + if usage_floor_startup == PersistedUsageFloorStartup::BootstrapPending + && !scanner_may_resume_pristine_usage_bootstrap(&cycle_info) + { + 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 = "usage_floor_bootstrap_conflict", + next_cycle = cycle_info.next, + "Scanner stopped because a pristine usage bootstrap conflicts with persisted cycle progress" + ); + global_metrics().set_cycle(None).await; + return Ok(()); + } + apply_persisted_usage_floor(&mut cycle_info, &mut leader_epoch, usage_floor); + let allow_pristine_bootstrap_pending = match usage_floor_startup { + PersistedUsageFloorStartup::Authoritative => false, + PersistedUsageFloorStartup::BootstrapPending => true, + PersistedUsageFloorStartup::Missing => { + if !may_bootstrap_missing_usage_floor || ctx.is_cancelled() || guard.is_lock_lost() { + global_metrics().set_cycle(None).await; + return Ok(()); + } + + let bootstrap_ctx = ctx.child_token(); + match await_scanner_cycle_with_lock_fence( + &bootstrap_ctx, + initialize_pristine_usage_baseline(storeapi.clone()), + guard.lock_lost_notified(), + ) + .await + { + Some(Ok(())) => true, + Some(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 = "usage_floor_bootstrap_failed", + error = %err, + "Scanner stopped because the pristine usage bootstrap could not be initialized" + ); + global_metrics().set_cycle(None).await; + return Ok(()); + } + None => { + global_metrics().set_cycle(None).await; + return Ok(()); + } + } } }; - apply_persisted_usage_floor(&mut cycle_info, &mut leader_epoch, usage_floor); if ctx.is_cancelled() || guard.is_lock_lost() { global_metrics().set_cycle(None).await; @@ -2126,7 +2222,14 @@ async fn run_data_scanner_with_maintenance_state( let claim_ctx = ctx.child_token(); let leadership_claimed = await_scanner_cycle_with_lock_fence( &claim_ctx, - claim_scanner_leadership(&claim_ctx, storeapi.clone(), &mut cycle_info, &mut cycle_revision, &mut leader_epoch), + claim_scanner_leadership( + &claim_ctx, + storeapi.clone(), + &mut cycle_info, + &mut cycle_revision, + &mut leader_epoch, + allow_pristine_bootstrap_pending, + ), guard.lock_lost_notified(), ) .await diff --git a/crates/scanner/src/scanner/cycle_state.rs b/crates/scanner/src/scanner/cycle_state.rs index 03a6b03ec..c09e8e360 100644 --- a/crates/scanner/src/scanner/cycle_state.rs +++ b/crates/scanner/src/scanner/cycle_state.rs @@ -926,7 +926,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< "scanner leader lock was lost after fencing newer cycle state".to_string(), )); } - fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), fence_epoch, Some(reset_epoch)) + fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), fence_epoch, Some(reset_epoch), false) .await .map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?; if guard.is_lock_lost() { @@ -1032,7 +1032,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc< "scanner leader lock was lost after rebuilding cycle state".to_string(), )); } - if let Err(err) = fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch)).await + if let Err(err) = + fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false).await { set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { path: DATA_USAGE_BLOOM_NAME_PATH.clone(), @@ -1179,6 +1180,13 @@ pub(super) struct PersistedUsageFloor { pub(super) leader_epoch: u64, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(super) enum PersistedUsageFloorStartup { + Authoritative, + Missing, + BootstrapPending, +} + pub(super) fn encode_scanner_cycle_state( cycle_info: &CurrentCycle, leader_epoch: u64, @@ -1300,11 +1308,25 @@ pub(super) fn advance_scanner_cycle(cycle_info: &mut CurrentCycle) -> Result<(), pub(super) async fn persisted_usage_floor( storeapi: Arc, ) -> Result { + let (floor, state) = persisted_usage_floor_for_startup(storeapi, false).await?; + if state != PersistedUsageFloorStartup::Authoritative { + return Err(ScannerError::Other( + "persisted scanner usage floor has no authoritative baseline".to_string(), + )); + } + Ok(floor) +} + +pub(super) async fn persisted_usage_floor_for_startup( + storeapi: Arc, + allow_missing_for_pristine_startup: bool, +) -> Result<(PersistedUsageFloor, PersistedUsageFloorStartup), ScannerError> { let Some(read_epoch) = scanner_publication_epoch(storeapi.clone()).await else { return Err(ScannerError::Other("scanner usage floor read is blocked by data movement".to_string())); }; let mut floor = PersistedUsageFloor::default(); let mut found_any = false; + let mut bootstrap_pending = false; let update_floor = |floor: &mut PersistedUsageFloor, usage: &DataUsageInfo, path: &str| -> Result<(), ScannerError> { floor.leader_epoch = floor.leader_epoch.max(usage.scanner_epoch.unwrap_or_default()); if let Some(completed_cycle) = usage.scanner_cycle { @@ -1323,14 +1345,24 @@ pub(super) async fn persisted_usage_floor( let usage = serde_json::from_slice::(&data).map_err(|err| { ScannerError::Other(format!("failed to decode scanner usage floor from {primary_path}: {err}")) })?; - if !data_usage_info_has_persisted_baseline_identity(&usage) { + if data_usage_info_is_pristine_bootstrap_pending(&usage) && primary_path == DATA_USAGE_OBJ_NAME_PATH.as_str() { + if bootstrap_pending { + return Err(ScannerError::Other( + "multiple pristine scanner usage bootstrap markers were found".to_string(), + )); + } + bootstrap_pending = true; + update_floor(&mut floor, &usage, primary_path)?; + None + } else if !data_usage_info_has_persisted_baseline_identity(&usage) { return Err(ScannerError::Other(format!( "scanner usage floor from {primary_path} has no persisted baseline identity" ))); + } else { + let epoch = usage.scanner_epoch.unwrap_or_default(); + update_floor(&mut floor, &usage, primary_path)?; + Some(epoch) } - let epoch = usage.scanner_epoch.unwrap_or_default(); - update_floor(&mut floor, &usage, primary_path)?; - Some(epoch) } Ok((None, _)) => None, Err(err) => { @@ -1342,6 +1374,11 @@ pub(super) async fn persisted_usage_floor( let mut any_found = primary_epoch.is_some(); match read_config_with_revision(storeapi.clone(), &backup_path).await { Ok((Some(data), _)) => { + if bootstrap_pending { + return Err(ScannerError::Other( + "pristine scanner usage bootstrap conflicts with a persisted backup".to_string(), + )); + } 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}")) @@ -1367,12 +1404,22 @@ pub(super) async fn persisted_usage_floor( } } if any_found { + if bootstrap_pending { + return Err(ScannerError::Other( + "pristine scanner usage bootstrap conflicts with an authoritative usage floor".to_string(), + )); + } found_any = true; break; } } - if !found_any { + if !found_any && !bootstrap_pending { + if !allow_missing_for_pristine_startup { + 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(), @@ -1405,7 +1452,14 @@ pub(super) async fn persisted_usage_floor( "scanner usage floor changed while its epoch proof was being confirmed".to_string(), )); }; - Ok(floor) + let state = if found_any { + PersistedUsageFloorStartup::Authoritative + } else if bootstrap_pending { + PersistedUsageFloorStartup::BootstrapPending + } else { + PersistedUsageFloorStartup::Missing + }; + Ok((floor, state)) } pub(super) fn apply_persisted_usage_floor(cycle_info: &mut CurrentCycle, leader_epoch: &mut u64, floor: PersistedUsageFloor) { diff --git a/crates/scanner/src/scanner/leadership.rs b/crates/scanner/src/scanner/leadership.rs index 69d54f2aa..80df0ff18 100644 --- a/crates/scanner/src/scanner/leadership.rs +++ b/crates/scanner/src/scanner/leadership.rs @@ -60,10 +60,18 @@ pub(super) async fn reconcile_scanner_leadership_claim( }) } -pub(super) fn decode_usage_snapshot_for_epoch_fence(data: &[u8], path: &str) -> Result { +pub(super) fn decode_usage_snapshot_for_epoch_fence( + data: &[u8], + path: &str, + allow_pristine_bootstrap_pending: bool, +) -> Result { let usage: DataUsageInfo = serde_json::from_slice(data) .map_err(|err| ScannerError::Other(format!("failed to decode scanner usage epoch fence from {path}: {err}")))?; - if !data_usage_info_has_persisted_baseline_identity(&usage) { + if !data_usage_info_has_persisted_baseline_identity(&usage) + && !(allow_pristine_bootstrap_pending + && path == DATA_USAGE_OBJ_NAME_PATH.as_str() + && data_usage_info_is_pristine_bootstrap_pending(&usage)) + { return Err(ScannerError::Other(format!( "scanner usage epoch fence from {path} has no persisted baseline identity" ))); @@ -74,9 +82,15 @@ pub(super) fn decode_usage_snapshot_for_epoch_fence(data: &[u8], path: &str) -> pub(super) async fn usage_snapshot_for_epoch_fence( storeapi: Arc, primary: Option<&[u8]>, + allow_pristine_bootstrap_pending: bool, ) -> Result, ScannerError> { if let Some(primary) = primary { - return decode_usage_snapshot_for_epoch_fence(primary, DATA_USAGE_OBJ_NAME_PATH.as_str()).map(Some); + return decode_usage_snapshot_for_epoch_fence( + primary, + DATA_USAGE_OBJ_NAME_PATH.as_str(), + allow_pristine_bootstrap_pending, + ) + .map(Some); } let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -84,7 +98,7 @@ pub(super) async fn usage_snapshot_for_epoch_fence( .await .map_err(|err| ScannerError::Other(format!("failed to read scanner usage epoch fence backup: {err}")))?; if let Some(backup) = backup.as_deref() { - return decode_usage_snapshot_for_epoch_fence(backup, &backup_path).map(Some); + return decode_usage_snapshot_for_epoch_fence(backup, &backup_path, false).map(Some); } for path in [ @@ -95,7 +109,7 @@ pub(super) async fn usage_snapshot_for_epoch_fence( .await .map_err(|err| ScannerError::Other(format!("failed to read legacy scanner usage epoch fence: {err}")))?; if let Some(legacy) = legacy.as_deref() { - return decode_usage_snapshot_for_epoch_fence(legacy, &path).map(Some); + return decode_usage_snapshot_for_epoch_fence(legacy, &path, false).map(Some); } } // A missing usage snapshot is an uninitialized state, not an empty @@ -104,36 +118,50 @@ 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( +pub(super) async fn initialize_pristine_usage_baseline( storeapi: Arc, - expected_publication_epoch: u64, -) -> Result { - let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_publication_epoch).await - else { +) -> Result<(), ScannerError> { + let Some(expected_epoch) = scanner_publication_epoch(storeapi.clone()).await else { return Err(ScannerError::Other( - "scanner publication epoch changed before confirming pristine usage state".to_string(), + "pristine scanner usage baseline initialization is blocked by data movement".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(), - )); + let baseline = DataUsageInfo { + last_update: Some(std::time::SystemTime::now()), + usage_snapshot_converged: Some(false), + usage_snapshot_bootstrap_pending: true, + ..Default::default() + }; + let data = serde_json::to_vec(&baseline) + .map_err(|err| ScannerError::Other(format!("failed to encode pristine scanner usage baseline: {err}")))?; + let save_result = save_config_with_publication_admission_for_epoch( + storeapi.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + data.clone(), + DataUsageCacheRevision::Missing.preconditions(), + expected_epoch, + ) + .await; + if save_result + .as_ref() + .ok() + .and_then(|info| info.etag.as_deref()) + .is_some_and(|etag| !etag.is_empty()) + { + return Ok(()); } - drop(publication_admission); - scanner_publication_admission_for_epoch(storeapi, expected_publication_epoch) + + let (persisted, revision) = read_config_with_revision(storeapi, DATA_USAGE_OBJ_NAME_PATH.as_str()) .await - .ok_or_else(|| ScannerError::Other("scanner publication epoch changed during pristine usage confirmation".to_string())) + .map_err(|err| ScannerError::Other(format!("failed to reconcile pristine scanner usage bootstrap: {err}")))?; + if persisted.as_deref() == Some(data.as_slice()) && matches!(revision, DataUsageCacheRevision::Etag(_)) { + return Ok(()); + } + + Err(ScannerError::Other(match save_result { + Ok(_) => "pristine scanner usage bootstrap returned no ETag and could not be confirmed".to_string(), + Err(err) => format!("failed to persist pristine scanner usage bootstrap: {err}"), + })) } pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( @@ -141,6 +169,7 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( storeapi: Arc, claimed_epoch: u64, expected_publication_epoch: Option, + allow_pristine_bootstrap_pending: bool, ) -> Result<(), ScannerError> { for retry in 0..=SCANNER_PERSIST_CAS_RETRIES { if ctx.is_cancelled() { @@ -160,10 +189,21 @@ 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 (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(()); + 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(), allow_pristine_bootstrap_pending).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())); }; match usage.scanner_epoch { Some(epoch) if epoch > claimed_epoch => { @@ -203,7 +243,11 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( .await .map_err(|err| ScannerError::Other(format!("failed to reconcile scanner usage epoch fence: {err}")))?; if let Some(persisted) = persisted { - let persisted = decode_usage_snapshot_for_epoch_fence(&persisted, DATA_USAGE_OBJ_NAME_PATH.as_str())?; + let persisted = decode_usage_snapshot_for_epoch_fence( + &persisted, + DATA_USAGE_OBJ_NAME_PATH.as_str(), + allow_pristine_bootstrap_pending, + )?; match persisted.scanner_epoch { Some(epoch) if epoch == claimed_epoch => return Ok(()), Some(epoch) if epoch > claimed_epoch => { @@ -233,9 +277,16 @@ pub(super) async fn complete_scanner_leadership_claim( storeapi: Arc, claimed_epoch: u64, expected_publication_epoch: Option, + allow_pristine_bootstrap_pending: bool, ) -> bool { - if let Err(err) = - fence_scanner_usage_epoch_with_expected_epoch(ctx, storeapi, claimed_epoch, expected_publication_epoch).await + if let Err(err) = fence_scanner_usage_epoch_with_expected_epoch( + ctx, + storeapi, + claimed_epoch, + expected_publication_epoch, + allow_pristine_bootstrap_pending, + ) + .await { error!( target: "rustfs::scanner", @@ -259,6 +310,7 @@ pub(super) async fn claim_scanner_leadership( cycle_info: &mut CurrentCycle, revision: &mut DataUsageCacheRevision, persisted_epoch: &mut u64, + allow_pristine_bootstrap_pending: bool, ) -> bool { for retry in 0..=SCANNER_PERSIST_CAS_RETRIES { if ctx.is_cancelled() { @@ -297,9 +349,8 @@ pub(super) async fn claim_scanner_leadership( let Some(read_epoch) = scanner_publication_epoch(storeapi.clone()).await else { return false; }; - let usage_baseline_missing = match read_usage_snapshot_for_epoch_fence(storeapi.clone()).await { - Ok((Some(_), _)) => false, - Ok((None, _)) => true, + let (usage_primary, _) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await { + Ok(result) => result, Err(err) => { error!( target: "rustfs::scanner", @@ -314,33 +365,40 @@ pub(super) async fn claim_scanner_leadership( return false; } }; + match usage_snapshot_for_epoch_fence(storeapi.clone(), usage_primary.as_deref(), allow_pristine_bootstrap_pending).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 _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; - } + let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { + if retry < SCANNER_PERSIST_CAS_RETRIES { + continue; } - } 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 + return false; }; save_config_with_preconditions(storeapi.clone(), &DATA_USAGE_BLOOM_NAME_PATH, data.clone(), revision.preconditions()) .await @@ -350,7 +408,14 @@ pub(super) async fn claim_scanner_leadership( if let Some(etag) = object_info.etag.filter(|etag| !etag.is_empty()) { *revision = DataUsageCacheRevision::Etag(etag); *persisted_epoch = claimed_epoch; - return complete_scanner_leadership_claim(ctx, storeapi, claimed_epoch, Some(read_epoch)).await; + return complete_scanner_leadership_claim( + ctx, + storeapi, + claimed_epoch, + Some(read_epoch), + allow_pristine_bootstrap_pending, + ) + .await; } match reconcile_scanner_leadership_claim( @@ -365,7 +430,14 @@ pub(super) async fn claim_scanner_leadership( .await { Ok(ScannerLeadershipClaimReconcile::Durable) => { - return complete_scanner_leadership_claim(ctx, storeapi, claimed_epoch, Some(read_epoch)).await; + return complete_scanner_leadership_claim( + ctx, + storeapi, + claimed_epoch, + Some(read_epoch), + allow_pristine_bootstrap_pending, + ) + .await; } Ok(ScannerLeadershipClaimReconcile::Changed) if retry < SCANNER_PERSIST_CAS_RETRIES => continue, Ok(ScannerLeadershipClaimReconcile::Changed | ScannerLeadershipClaimReconcile::Unchanged) => { @@ -409,7 +481,14 @@ pub(super) async fn claim_scanner_leadership( .await { Ok(ScannerLeadershipClaimReconcile::Durable) => { - return complete_scanner_leadership_claim(ctx, storeapi, claimed_epoch, Some(read_epoch)).await; + return complete_scanner_leadership_claim( + ctx, + storeapi, + claimed_epoch, + Some(read_epoch), + allow_pristine_bootstrap_pending, + ) + .await; } Ok(ScannerLeadershipClaimReconcile::Changed) if retry < SCANNER_PERSIST_CAS_RETRIES && !ctx.is_cancelled() => diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 3ed098146..58b919e44 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -373,18 +373,6 @@ async fn insert_usage_after_first_legacy_backup_read(store: &MemoryConfigStore) ); } -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; @@ -2090,6 +2078,21 @@ async fn scanner_usage_floor_fails_closed_on_corrupt_or_exhausted_usage_state() "a structurally incomplete usage snapshot must not be treated as an empty floor" ); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), + serde_json::to_vec(&DataUsageInfo { + last_update: Some(std::time::SystemTime::now()), + scanner_cycle: Some(1), + usage_snapshot_bootstrap_pending: true, + ..Default::default() + }) + .expect("pending usage marker should encode"), + ); + assert!( + persisted_usage_floor(store.clone()).await.is_err(), + "a pending marker must never pass the legacy authoritative fallback" + ); + store.objects.lock().await.insert( memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), serde_json::to_vec(&DataUsageInfo { @@ -2101,6 +2104,26 @@ 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_allows_only_explicit_pristine_bootstrap() { + let store = Arc::new(MemoryConfigStore::default()); + let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("a verified pristine startup should use the empty floor"); + assert_eq!(floor, PersistedUsageFloor::default()); + assert_eq!(state, PersistedUsageFloorStartup::Missing); + assert!(persisted_usage_floor_for_startup(store.clone(), false).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(), + ); + assert!( + persisted_usage_floor_for_startup(store, true).await.is_err(), + "pristine bootstrap must not hide corrupt persisted state" + ); +} + #[tokio::test] async fn scanner_usage_floor_fails_closed_on_zero_byte_usage_objects() { for path in [ @@ -2162,6 +2185,86 @@ async fn scanner_usage_floor_rejects_publication_change_during_pristine_confirma assert!(persisted_usage_floor(store).await.is_err()); } +#[test] +fn scanner_pristine_cycle_state_requires_no_durable_progress() { + let cycle = CurrentCycle::default(); + assert!(scanner_cycle_state_is_pristine(&cycle, 0, &DataUsageCacheRevision::Missing)); + assert!(scanner_may_resume_pristine_usage_bootstrap(&cycle)); + assert!(!scanner_cycle_state_is_pristine( + &CurrentCycle { + next: 1, + ..Default::default() + }, + 0, + &DataUsageCacheRevision::Missing + )); + assert!(!scanner_may_resume_pristine_usage_bootstrap(&CurrentCycle { + next: 1, + ..Default::default() + })); + assert!(!scanner_cycle_state_is_pristine(&cycle, 1, &DataUsageCacheRevision::Missing)); + assert!(!scanner_cycle_state_is_pristine( + &cycle, + 0, + &DataUsageCacheRevision::Etag("etag".to_string()) + )); +} + +#[tokio::test] +async fn scanner_pristine_bootstrap_allows_first_bucket_to_win_startup() { + let (_temp_dir, store) = setup_scanner_cycle_store().await; + let cycle = CurrentCycle::default(); + let revision = DataUsageCacheRevision::Missing; + + store + .make_bucket("first-user-bucket", &crate::storage_api::scan::MakeBucketOptions::default()) + .await + .expect("test bucket should be created"); + + assert!(scanner_may_bootstrap_missing_usage_floor(&cycle, 0, &revision)); + assert_eq!( + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("first startup should still admit a non-authoritative bootstrap marker") + .1, + PersistedUsageFloorStartup::Missing + ); + initialize_pristine_usage_baseline(store.clone()) + .await + .expect("first startup should persist its pending marker"); + let pending = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("pending marker should be stored"); + let pending = serde_json::from_slice::(&pending).expect("pending marker should decode"); + assert!(data_usage_info_is_pristine_bootstrap_pending(&pending)); + assert!(!data_usage_info_has_persisted_baseline_identity(&pending)); + assert_eq!( + persisted_usage_floor_for_startup(store.clone(), false) + .await + .expect("the pending marker should be resumable after restart") + .1, + PersistedUsageFloorStartup::BootstrapPending + ); + + store + .delete_bucket("first-user-bucket", &crate::storage_api::scan::DeleteBucketOptions::default()) + .await + .expect("first user bucket should be deleted"); + assert!( + read_config(store.clone(), &format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str())) + .await + .is_err(), + "bucket deletion must not copy the pending marker into the backup slot" + ); + assert_eq!( + persisted_usage_floor_for_startup(store, false) + .await + .expect("the pending marker should remain resumable after bucket deletion") + .1, + PersistedUsageFloorStartup::BootstrapPending + ); +} + #[tokio::test] async fn scanner_usage_backup_uses_durable_cycle_cadence_across_tasks() { let store = Arc::new(MemoryConfigStore::default()); @@ -2444,7 +2547,7 @@ async fn test_leadership_claim_preserves_usage_epoch_floor_across_old_epoch_conf ); let mut persisted_epoch = 8; - assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch,).await); + assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await); let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) .await @@ -2467,50 +2570,111 @@ async fn test_leadership_claim_rejects_terminal_epoch() { }; let mut persisted_epoch = u64::MAX - 1; - assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch).await); + assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await); assert_eq!(persisted_epoch, u64::MAX - 1); assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); } #[tokio::test] -async fn scanner_bootstraps_leadership_when_usage_snapshots_are_stably_absent() { +async fn scanner_defers_leadership_when_usage_snapshots_are_stably_absent() { let store = Arc::new(MemoryConfigStore::default()); - 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); + 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, false).await); + assert!(read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH).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() { +async fn pristine_usage_bootstrap_pending_unblocks_first_leadership_claim() { let store = Arc::new(MemoryConfigStore::default()); - insert_usage_after_first_legacy_backup_read(store.as_ref()).await; + initialize_pristine_usage_baseline(store.clone()) + .await + .expect("verified pristine startup should publish its pending marker"); - 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()); + let ctx = CancellationToken::new(); + let mut revision = DataUsageCacheRevision::Missing; + let mut cycle = CurrentCycle::default(); + let mut persisted_epoch = 0; + assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false,).await); + assert!(read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); + assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, true).await); + + let usage = read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("leadership claim should fence the pristine baseline"); + let usage = serde_json::from_slice::(&usage).expect("pristine bootstrap marker should remain valid"); + assert!(data_usage_info_is_pristine_bootstrap_pending(&usage)); + assert!(!data_usage_info_has_persisted_baseline_identity(&usage)); + assert_eq!(usage.scanner_epoch, Some(1)); } #[tokio::test] -async fn leadership_claim_rejects_publication_change_during_pristine_confirmation() { +async fn existing_pristine_usage_bootstrap_is_resumed_after_restart() { let store = Arc::new(MemoryConfigStore::default()); - store.block_publication_after_admissions.store(2, Ordering::Release); + initialize_pristine_usage_baseline(store.clone()) + .await + .expect("verified pristine startup should publish its pending marker"); - 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()); + let (floor, state) = persisted_usage_floor_for_startup(store.clone(), false) + .await + .expect("restart should recognize the pending pristine bootstrap"); + assert_eq!(floor, PersistedUsageFloor::default()); + assert_eq!(state, PersistedUsageFloorStartup::BootstrapPending); + assert!(persisted_usage_floor(store).await.is_err()); +} + +#[tokio::test] +async fn pristine_usage_bootstrap_reconciles_post_commit_error() { + let store = Arc::new(MemoryConfigStore::default()); + let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store.error_after_commit_put_number.lock().await.insert(key, 1); + + initialize_pristine_usage_baseline(store.clone()) + .await + .expect("a committed pending marker should reconcile after a lost response"); + let usage = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("the reconciled pending marker should remain"); + let usage = serde_json::from_slice::(&usage).expect("pending marker should decode"); + assert!(data_usage_info_is_pristine_bootstrap_pending(&usage)); + assert!(!data_usage_info_has_persisted_baseline_identity(&usage)); + assert_eq!( + persisted_usage_floor_for_startup(store.clone(), false) + .await + .expect("restart should resume a committed pending marker") + .1, + PersistedUsageFloorStartup::BootstrapPending + ); + assert!(persisted_usage_floor(store).await.is_err()); +} + +#[tokio::test] +async fn pristine_usage_bootstrap_does_not_overwrite_concurrent_replacement() { + let store = Arc::new(MemoryConfigStore::default()); + let key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + let replacement = serde_json::to_vec(&complete_usage_with_bucket_count(None, 1)).expect("replacement should encode"); + store + .replace_after_successful_puts + .lock() + .await + .insert(key, (1, replacement.clone())); + + initialize_pristine_usage_baseline(store.clone()) + .await + .expect("the bootstrap write completed before the replacement"); + assert_eq!( + read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("newer usage snapshot must remain"), + replacement + ); } #[tokio::test] @@ -2528,7 +2692,7 @@ async fn leadership_claim_defers_on_corrupt_usage_baseline_without_bloom_write() }; let mut persisted_epoch = 0; - assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch,).await); + assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await); assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); } @@ -2548,7 +2712,7 @@ async fn leadership_claim_defers_on_unidentified_usage_baseline_without_bloom_wr }; let mut persisted_epoch = 0; - assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch,).await); + assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await); assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); } @@ -2570,7 +2734,7 @@ async fn test_leadership_claim_confirms_commit_after_returned_error() { let mut persisted_epoch = 0; seed_usage_snapshot_for_leadership_claim(&store).await; - assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch).await); + assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await); let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) .await @@ -2626,7 +2790,7 @@ async fn test_leadership_claim_usage_fence_rejects_old_inflight_writer() { ..Default::default() }; let mut persisted_epoch = 4; - assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch).await); + assert!(claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await); let (fenced_data, fenced_revision) = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) .await @@ -2691,6 +2855,7 @@ async fn cycle_budget_lease_takeover_rejects_old_generation() { &mut replacement_cycle, &mut replacement_revision, &mut replacement_epoch, + false, ) .await );