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
This commit is contained in:
Zhengchao An
2026-08-24 14:29:09 +08:00
committed by GitHub
parent fc98dbb654
commit 170a4c7640
7 changed files with 582 additions and 144 deletions
+8 -2
View File
@@ -573,6 +573,10 @@ pub struct DataUsageInfo {
/// cycle (or retained a compatible last-known-good cache). /// cycle (or retained a compatible last-known-good cache).
#[serde(default)] #[serde(default)]
pub usage_snapshot_partial: bool, 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 /// Deprecated kept here for backward compatibility reasons
pub bucket_sizes: HashMap<String, u64>, pub bucket_sizes: HashMap<String, u64>,
/// Per-disk snapshot information when available /// Per-disk snapshot information when available
@@ -1962,7 +1966,8 @@ impl DataUsageInfo {
/// Whether this snapshot authoritatively covers every reported bucket. /// Whether this snapshot authoritatively covers every reported bucket.
pub fn is_complete_bucket_usage_snapshot(&self) -> bool { 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() && self.last_update.is_some()
&& u64::try_from(self.buckets_usage.len()).ok() == Some(self.buckets_count) && 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, /// admin display. Partial data is accepted only with unique set states,
/// a plan digest for every state, and at least one usable generation. /// a plan digest for every state, and at least one usable generation.
pub fn is_valid_partial_snapshot(&self) -> bool { 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.usage_snapshot_converged != Some(false)
|| self.last_update.is_none() || self.last_update.is_none()
|| self.scanner_cycle.is_none() || self.scanner_cycle.is_none()
+3 -3
View File
@@ -33,7 +33,7 @@ use crate::bucket::{
use crate::cache_value::metacache_set::{ListPathRawOptions, list_path_raw}; use crate::cache_value::metacache_set::{ListPathRawOptions, list_path_raw};
use crate::config::com::{ use crate::config::com::{
CONFIG_PREFIX, delete_config, read_config_limited_preserve_empty, read_config_limited_preserve_empty_with_metadata, 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, save_config_with_opts_quiet,
}; };
use crate::data_movement; use crate::data_movement;
@@ -4848,7 +4848,7 @@ impl ECStore {
let mut pool_meta = self.pool_meta.write().await; let mut pool_meta = self.pool_meta.write().await;
record_decommission_unresolved_entry(&mut pool_meta, idx, generation, entry)?; record_decommission_unresolved_entry(&mut pool_meta, idx, generation, entry)?;
} }
self.save_current_pool_meta() self.save_current_pool_meta(&[idx])
.await .await
.map_err(|err| Error::other(format!("decommission unresolved entry ledger save failed: {err}"))) .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.run_guarded_decommission_side_effect(rx, &operation_gate, || {
self.check_after_decommission_unfenced(idx, generation) self.check_after_decommission_unfenced(idx, generation)
}) })
.await .await
} }
async fn check_after_decommission_unfenced( async fn check_after_decommission_unfenced(
+32 -1
View File
@@ -743,7 +743,12 @@ where
let backup_seed = load_data_usage_for_bucket_removal(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) let backup_seed = load_data_usage_for_bucket_removal(store, DATA_USAGE_OBJ_NAME_PATH.as_str())
.await? .await?
.map(|(data_usage_info, _)| data_usage_info) .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( remove_bucket_usage_from_object_with_retries_and_publication(
store, store,
DATA_USAGE_OBJ_BACKUP_PATH.as_str(), DATA_USAGE_OBJ_BACKUP_PATH.as_str(),
@@ -4646,6 +4651,32 @@ mod tests {
assert!(state.backup_object.is_none()); 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::<DataUsageInfo>(&state.object.as_ref().expect("pending primary should remain").0)
.expect("pending primary should decode");
assert!(saved.usage_snapshot_bootstrap_pending);
}
#[tokio::test] #[tokio::test]
async fn remove_bucket_usage_migrates_legacy_snapshot_without_hiding_other_buckets() { 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()); let mut legacy = data_usage_info_for_test("bucket-a", 2, 84, SystemTime::now());
+122 -19
View File
@@ -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. // complete: a timestamp, a scanner cycle, and an exact bucket cardinality.
// A current snapshot with only scanner_epoch/scanner_cycle (or an explicit // A current snapshot with only scanner_epoch/scanner_cycle (or an explicit
// incomplete marker) is not evidence of a durable usage baseline. // 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.scanner_epoch.is_none()
&& info.usage_snapshot_converged != Some(false) && info.usage_snapshot_converged != Some(false)
&& info.last_update.is_some() && 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) && 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 { fn usage_cache_needs_prompt_scan(authoritative: &DataUsageInfo, observed: Option<&DataUsageInfo>) -> bool {
data_usage_info_is_cold(authoritative) data_usage_info_is_cold(authoritative)
|| observed.is_some_and(|observed| observed_data_usage_is_newer(observed, 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<ECStore>) -> (bool,
(persisted_usage_cache_is_cold_for_startup(storeapi).await, has_buckets) (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<ECStore>) { pub async fn init_data_scanner(ctx: CancellationToken, storeapi: Arc<ECStore>) {
let (startup_features, startup_maintenance_generation) = configure_scanner_defaults(&ctx, &storeapi).await; let (startup_features, startup_maintenance_generation) = configure_scanner_defaults(&ctx, &storeapi).await;
// Force init global sleeper so config is read once at startup. // Force init global sleeper so config is read once at startup.
@@ -1153,7 +1192,7 @@ where
LockLost: Future<Output = ()>, LockLost: Future<Output = ()>,
{ {
let fence_ctx = ctx.child_token(); 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!(claim);
tokio::pin!(lock_lost); tokio::pin!(lock_lost);
tokio::select! { tokio::select! {
@@ -2100,24 +2139,81 @@ async fn run_data_scanner_with_maintenance_state(
return Err(err); return Err(err);
} }
}; };
let usage_floor = match persisted_usage_floor(storeapi.clone()).await { let may_bootstrap_missing_usage_floor = scanner_may_bootstrap_missing_usage_floor(&cycle_info, leader_epoch, &cycle_revision);
Ok(floor) => floor, let (usage_floor, usage_floor_startup) =
Err(err) => { match persisted_usage_floor_for_startup(storeapi.clone(), may_bootstrap_missing_usage_floor).await {
error!( Ok(result) => result,
target: "rustfs::scanner", Err(err) => {
event = EVENT_SCANNER_PERSIST_STATE, error!(
component = LOG_COMPONENT_SCANNER, target: "rustfs::scanner",
subsystem = LOG_SUBSYSTEM_RUNTIME, event = EVENT_SCANNER_PERSIST_STATE,
path = %DATA_USAGE_OBJ_NAME_PATH.as_str(), component = LOG_COMPONENT_SCANNER,
state = "usage_floor_load_failed", subsystem = LOG_SUBSYSTEM_RUNTIME,
error = %err, path = %DATA_USAGE_OBJ_NAME_PATH.as_str(),
"Scanner stopped because the persisted usage floor could not be loaded" state = "usage_floor_load_failed",
); error = %err,
global_metrics().set_cycle(None).await; "Scanner stopped because the persisted usage floor could not be loaded"
return Ok(()); );
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() { if ctx.is_cancelled() || guard.is_lock_lost() {
global_metrics().set_cycle(None).await; 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 claim_ctx = ctx.child_token();
let leadership_claimed = await_scanner_cycle_with_lock_fence( let leadership_claimed = await_scanner_cycle_with_lock_fence(
&claim_ctx, &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(), guard.lock_lost_notified(),
) )
.await .await
+62 -8
View File
@@ -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(), "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 .await
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?; .map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
if guard.is_lock_lost() { 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(), "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 { set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus {
path: DATA_USAGE_BLOOM_NAME_PATH.clone(), path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
@@ -1179,6 +1180,13 @@ pub(super) struct PersistedUsageFloor {
pub(super) leader_epoch: u64, 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( pub(super) fn encode_scanner_cycle_state(
cycle_info: &CurrentCycle, cycle_info: &CurrentCycle,
leader_epoch: u64, 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( pub(super) async fn persisted_usage_floor(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>, storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
) -> Result<PersistedUsageFloor, ScannerError> { ) -> Result<PersistedUsageFloor, ScannerError> {
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<impl ScannerObjectIO + ScannerConfigObjectDelete>,
allow_missing_for_pristine_startup: bool,
) -> Result<(PersistedUsageFloor, PersistedUsageFloorStartup), ScannerError> {
let Some(read_epoch) = scanner_publication_epoch(storeapi.clone()).await else { 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())); return Err(ScannerError::Other("scanner usage floor read is blocked by data movement".to_string()));
}; };
let mut floor = PersistedUsageFloor::default(); let mut floor = PersistedUsageFloor::default();
let mut found_any = false; let mut found_any = false;
let mut bootstrap_pending = false;
let update_floor = |floor: &mut PersistedUsageFloor, usage: &DataUsageInfo, path: &str| -> Result<(), ScannerError> { 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()); floor.leader_epoch = floor.leader_epoch.max(usage.scanner_epoch.unwrap_or_default());
if let Some(completed_cycle) = usage.scanner_cycle { 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::<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}"))
})?; })?;
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!( return Err(ScannerError::Other(format!(
"scanner usage floor from {primary_path} has no persisted baseline identity" "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, Ok((None, _)) => None,
Err(err) => { Err(err) => {
@@ -1342,6 +1374,11 @@ 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_with_revision(storeapi.clone(), &backup_path).await { match read_config_with_revision(storeapi.clone(), &backup_path).await {
Ok((Some(data), _)) => { Ok((Some(data), _)) => {
if bootstrap_pending {
return Err(ScannerError::Other(
"pristine scanner usage bootstrap conflicts with a persisted backup".to_string(),
));
}
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}"))
@@ -1367,12 +1404,22 @@ pub(super) async fn persisted_usage_floor(
} }
} }
if any_found { 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; found_any = true;
break; 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 { let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else {
return Err(ScannerError::Other( return Err(ScannerError::Other(
"scanner usage floor changed before pristine state confirmation".to_string(), "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(), "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) { pub(super) fn apply_persisted_usage_floor(cycle_info: &mut CurrentCycle, leader_epoch: &mut u64, floor: PersistedUsageFloor) {
+146 -67
View File
@@ -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<DataUsageInfo, ScannerError> { pub(super) fn decode_usage_snapshot_for_epoch_fence(
data: &[u8],
path: &str,
allow_pristine_bootstrap_pending: bool,
) -> Result<DataUsageInfo, ScannerError> {
let usage: DataUsageInfo = serde_json::from_slice(data) let usage: DataUsageInfo = serde_json::from_slice(data)
.map_err(|err| ScannerError::Other(format!("failed to decode scanner usage epoch fence from {path}: {err}")))?; .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!( return Err(ScannerError::Other(format!(
"scanner usage epoch fence from {path} has no persisted baseline identity" "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( pub(super) async fn usage_snapshot_for_epoch_fence(
storeapi: Arc<impl ScannerObjectIO>, storeapi: Arc<impl ScannerObjectIO>,
primary: Option<&[u8]>, primary: Option<&[u8]>,
allow_pristine_bootstrap_pending: bool,
) -> Result<Option<DataUsageInfo>, ScannerError> { ) -> Result<Option<DataUsageInfo>, ScannerError> {
if let Some(primary) = primary { 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()); 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 .await
.map_err(|err| ScannerError::Other(format!("failed to read scanner usage epoch fence backup: {err}")))?; .map_err(|err| ScannerError::Other(format!("failed to read scanner usage epoch fence backup: {err}")))?;
if let Some(backup) = backup.as_deref() { 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 [ for path in [
@@ -95,7 +109,7 @@ pub(super) async fn usage_snapshot_for_epoch_fence(
.await .await
.map_err(|err| ScannerError::Other(format!("failed to read legacy scanner usage epoch fence: {err}")))?; .map_err(|err| ScannerError::Other(format!("failed to read legacy scanner usage epoch fence: {err}")))?;
if let Some(legacy) = legacy.as_deref() { 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 // 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) Ok(None)
} }
async fn read_usage_snapshot_for_epoch_fence( pub(super) async fn initialize_pristine_usage_baseline(
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>, storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
expected_publication_epoch: u64, ) -> Result<(), ScannerError> {
) -> Result<crate::ScannerDataUsagePublicationAdmission, ScannerError> { let Some(expected_epoch) = scanner_publication_epoch(storeapi.clone()).await else {
let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_publication_epoch).await
else {
return Err(ScannerError::Other( 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?; let baseline = DataUsageInfo {
if usage.is_some() { last_update: Some(std::time::SystemTime::now()),
return Err(ScannerError::Other( usage_snapshot_converged: Some(false),
"scanner usage baseline appeared while confirming pristine bootstrap".to_string(), 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 .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( 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<impl ScannerObjectIO + ScannerConfigObjectDelete>, storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
claimed_epoch: u64, claimed_epoch: u64,
expected_publication_epoch: Option<u64>, expected_publication_epoch: Option<u64>,
allow_pristine_bootstrap_pending: bool,
) -> Result<(), ScannerError> { ) -> Result<(), ScannerError> {
for retry in 0..=SCANNER_PERSIST_CAS_RETRIES { for retry in 0..=SCANNER_PERSIST_CAS_RETRIES {
if ctx.is_cancelled() { 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(), "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 (primary, revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
let Some(mut usage) = usage else { .await
let _publication_admission = confirm_usage_snapshot_absent_for_bootstrap(storeapi, read_epoch).await?; .map_err(|err| ScannerError::Other(format!("failed to read scanner usage epoch fence: {err}")))?;
return Ok(()); 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 { match usage.scanner_epoch {
Some(epoch) if epoch > claimed_epoch => { Some(epoch) if epoch > claimed_epoch => {
@@ -203,7 +243,11 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
.await .await
.map_err(|err| ScannerError::Other(format!("failed to reconcile scanner usage epoch fence: {err}")))?; .map_err(|err| ScannerError::Other(format!("failed to reconcile scanner usage epoch fence: {err}")))?;
if let Some(persisted) = persisted { 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 { match persisted.scanner_epoch {
Some(epoch) if epoch == claimed_epoch => return Ok(()), Some(epoch) if epoch == claimed_epoch => return Ok(()),
Some(epoch) if epoch > claimed_epoch => { Some(epoch) if epoch > claimed_epoch => {
@@ -233,9 +277,16 @@ pub(super) async fn complete_scanner_leadership_claim(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>, storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
claimed_epoch: u64, claimed_epoch: u64,
expected_publication_epoch: Option<u64>, expected_publication_epoch: Option<u64>,
allow_pristine_bootstrap_pending: bool,
) -> bool { ) -> bool {
if let Err(err) = if let Err(err) = fence_scanner_usage_epoch_with_expected_epoch(
fence_scanner_usage_epoch_with_expected_epoch(ctx, storeapi, claimed_epoch, expected_publication_epoch).await ctx,
storeapi,
claimed_epoch,
expected_publication_epoch,
allow_pristine_bootstrap_pending,
)
.await
{ {
error!( error!(
target: "rustfs::scanner", target: "rustfs::scanner",
@@ -259,6 +310,7 @@ pub(super) async fn claim_scanner_leadership(
cycle_info: &mut CurrentCycle, cycle_info: &mut CurrentCycle,
revision: &mut DataUsageCacheRevision, revision: &mut DataUsageCacheRevision,
persisted_epoch: &mut u64, persisted_epoch: &mut u64,
allow_pristine_bootstrap_pending: bool,
) -> bool { ) -> bool {
for retry in 0..=SCANNER_PERSIST_CAS_RETRIES { for retry in 0..=SCANNER_PERSIST_CAS_RETRIES {
if ctx.is_cancelled() { 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 { let Some(read_epoch) = scanner_publication_epoch(storeapi.clone()).await else {
return false; return false;
}; };
let usage_baseline_missing = match read_usage_snapshot_for_epoch_fence(storeapi.clone()).await { let (usage_primary, _) = match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await {
Ok((Some(_), _)) => false, Ok(result) => result,
Ok((None, _)) => true,
Err(err) => { Err(err) => {
error!( error!(
target: "rustfs::scanner", target: "rustfs::scanner",
@@ -314,33 +365,40 @@ pub(super) async fn claim_scanner_leadership(
return false; 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 save_result = {
let _publication_admission = if usage_baseline_missing { let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else {
match confirm_usage_snapshot_absent_for_bootstrap(storeapi.clone(), read_epoch).await { if retry < SCANNER_PERSIST_CAS_RETRIES {
Ok(publication_admission) => publication_admission, continue;
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;
}
} }
} else { return false;
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
@@ -350,7 +408,14 @@ pub(super) async fn claim_scanner_leadership(
if let Some(etag) = object_info.etag.filter(|etag| !etag.is_empty()) { if let Some(etag) = object_info.etag.filter(|etag| !etag.is_empty()) {
*revision = DataUsageCacheRevision::Etag(etag); *revision = DataUsageCacheRevision::Etag(etag);
*persisted_epoch = claimed_epoch; *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( match reconcile_scanner_leadership_claim(
@@ -365,7 +430,14 @@ pub(super) async fn claim_scanner_leadership(
.await .await
{ {
Ok(ScannerLeadershipClaimReconcile::Durable) => { 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) if retry < SCANNER_PERSIST_CAS_RETRIES => continue,
Ok(ScannerLeadershipClaimReconcile::Changed | ScannerLeadershipClaimReconcile::Unchanged) => { Ok(ScannerLeadershipClaimReconcile::Changed | ScannerLeadershipClaimReconcile::Unchanged) => {
@@ -409,7 +481,14 @@ pub(super) async fn claim_scanner_leadership(
.await .await
{ {
Ok(ScannerLeadershipClaimReconcile::Durable) => { 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) Ok(ScannerLeadershipClaimReconcile::Changed)
if retry < SCANNER_PERSIST_CAS_RETRIES && !ctx.is_cancelled() => if retry < SCANNER_PERSIST_CAS_RETRIES && !ctx.is_cancelled() =>
+209 -44
View File
@@ -373,18 +373,6 @@ async fn insert_usage_after_first_legacy_backup_read(store: &MemoryConfigStore)
); );
} }
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;
@@ -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" "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( 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()),
serde_json::to_vec(&DataUsageInfo { 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()); 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] #[tokio::test]
async fn scanner_usage_floor_fails_closed_on_zero_byte_usage_objects() { async fn scanner_usage_floor_fails_closed_on_zero_byte_usage_objects() {
for path in [ 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()); 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::<DataUsageInfo>(&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] #[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());
@@ -2444,7 +2547,7 @@ async fn test_leadership_claim_preserves_usage_epoch_floor_across_old_epoch_conf
); );
let mut persisted_epoch = 8; 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) let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH)
.await .await
@@ -2467,50 +2570,111 @@ async fn test_leadership_claim_rejects_terminal_epoch() {
}; };
let mut persisted_epoch = u64::MAX - 1; 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_eq!(persisted_epoch, u64::MAX - 1);
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err());
} }
#[tokio::test] #[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 store = Arc::new(MemoryConfigStore::default());
let floor = persisted_usage_floor(store.clone()) let ctx = CancellationToken::new();
.await let mut revision = DataUsageCacheRevision::Missing;
.expect("pristine usage state should provide the initial floor"); let mut cycle = CurrentCycle {
assert_eq!(floor, PersistedUsageFloor::default()); next: 12,
let (claimed, persisted_epoch) = claim_test_scanner_leadership(store.clone()).await; ..Default::default()
assert!(claimed); };
let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) let mut persisted_epoch = 0;
.await
.expect("pristine leadership claim should persist"); assert!(!claim_scanner_leadership(&ctx, store.clone(), &mut cycle, &mut revision, &mut persisted_epoch, false).await);
let (persisted_cycle, leader_epoch) = decode_scanner_cycle_state(&state).expect("persisted leadership claim should decode"); assert!(read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH).await.is_err());
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] #[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()); 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; let ctx = CancellationToken::new();
assert!(!claimed); let mut revision = DataUsageCacheRevision::Missing;
assert_eq!(persisted_epoch, 0); let mut cycle = CurrentCycle::default();
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); 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::<DataUsageInfo>(&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] #[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()); 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; let (floor, state) = persisted_usage_floor_for_startup(store.clone(), false)
assert!(!claimed); .await
assert_eq!(persisted_epoch, 0); .expect("restart should recognize the pending pristine bootstrap");
assert!(read_config(store, &DATA_USAGE_BLOOM_NAME_PATH).await.is_err()); 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::<DataUsageInfo>(&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] #[tokio::test]
@@ -2528,7 +2692,7 @@ async fn leadership_claim_defers_on_corrupt_usage_baseline_without_bloom_write()
}; };
let mut persisted_epoch = 0; 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()); 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; 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()); 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; let mut persisted_epoch = 0;
seed_usage_snapshot_for_leadership_claim(&store).await; 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) let state = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH)
.await .await
@@ -2626,7 +2790,7 @@ async fn test_leadership_claim_usage_fence_rejects_old_inflight_writer() {
..Default::default() ..Default::default()
}; };
let mut persisted_epoch = 4; 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()) let (fenced_data, fenced_revision) = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await .await
@@ -2691,6 +2855,7 @@ async fn cycle_budget_lease_takeover_rejects_old_generation() {
&mut replacement_cycle, &mut replacement_cycle,
&mut replacement_revision, &mut replacement_revision,
&mut replacement_epoch, &mut replacement_epoch,
false,
) )
.await .await
); );