From fca1514aac8f82ab9c6f1114be7a36880a03d390 Mon Sep 17 00:00:00 2001 From: houseme Date: Mon, 31 Aug 2026 06:17:26 +0800 Subject: [PATCH] fix(scanner): recover legacy empty usage floor (#6914) Co-authored-by: heihutu --- crates/scanner/src/data_usage_define.rs | 5 + crates/scanner/src/scanner.rs | 68 ++- crates/scanner/src/scanner/cycle_state.rs | 474 +++++++++++++++- crates/scanner/src/scanner/tests.rs | 566 ++++++++++++++++++- docs/architecture/compat-cleanup-register.md | 1 + 5 files changed, 1094 insertions(+), 20 deletions(-) diff --git a/crates/scanner/src/data_usage_define.rs b/crates/scanner/src/data_usage_define.rs index 8447c4917..9d72f72f2 100644 --- a/crates/scanner/src/data_usage_define.rs +++ b/crates/scanner/src/data_usage_define.rs @@ -175,6 +175,11 @@ pub static DATA_USAGE_BUCKET: LazyLock = pub static DATA_USAGE_OBJ_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_OBJECT_NAME}")); +/// Durable evidence for recovery of the exact empty usage fence written by +/// rc.2/rc.3 bucket cleanup before the first authoritative scanner snapshot. +pub static DATA_USAGE_RECOVERY_PATH: LazyLock = + LazyLock::new(|| format!("{}.recovery-pending.json", DATA_USAGE_OBJ_NAME_PATH.as_str())); + pub static DATA_USAGE_OBSERVED_OBJ_NAME_PATH: LazyLock = LazyLock::new(|| format!("{BUCKET_META_PREFIX}{SLASH_SEPARATOR}{DATA_USAGE_OBSERVED_OBJECT_NAME}")); diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index a83b8a737..66d563362 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -22,7 +22,8 @@ use std::sync::{Arc, LazyLock, RwLock}; use self::heal_info::{BackgroundHealInfoReadStatus, read_background_heal_info_with_epoch, save_background_heal_info_for_epoch}; use crate::data_usage_define::{ BACKGROUND_HEAL_INFO_PATH, DATA_USAGE_BLOOM_NAME_PATH, DATA_USAGE_OBJ_NAME_PATH, DATA_USAGE_OBSERVED_OBJ_NAME_PATH, - DataUsageCache, DataUsageCacheRevision, LEGACY_DATA_USAGE_OBJ_NAME_PATH, read_config_revision, read_config_with_revision, + DATA_USAGE_RECOVERY_PATH, DataUsageCache, DataUsageCacheRevision, LEGACY_DATA_USAGE_OBJ_NAME_PATH, read_config_revision, + read_config_with_revision, }; use crate::runtime_config::{ ScannerRuntimeConfig, ScannerRuntimeConfigSource, refresh_scanner_runtime_config_from_global, scanner_bitrot_cycle, @@ -294,6 +295,14 @@ fn record_scanner_leader_lock_state(state: &'static str) { .increment(1); } +async fn finish_scanner_leader_iteration(lock_lost: bool, state: &'static str, error: String) { + reset_scanner_cycle_schedule(); + let liveness_already_recorded = lock_lost && !global_metrics().report().await.leader_lock_held_by_this_process; + if !liveness_already_recorded { + global_metrics().record_scanner_leader_liveness(state, false, error).await; + } +} + #[cfg(test)] fn scanner_cycle_max_duration() -> Option { resolve_scanner_runtime_config().cycle_budget.max_duration @@ -761,6 +770,13 @@ fn prepare_cycle_for_usage_floor_bootstrap( } (true, usage_floor.leader_epoch == 0) } + PersistedUsageFloorStartup::RecoveredLegacyEmptyFence => { + // The legacy empty fence proves only its leader epoch, not + // namespace coverage. Restart coverage from zero while retaining + // that epoch as the lower bound for the next leadership claim. + *cycle_info = CurrentCycle::default(); + (true, true) + } } } @@ -2240,6 +2256,7 @@ async fn run_data_scanner_with_maintenance_state( { let Some((features, generation)) = detect_stable_scanner_maintenance_features(&ctx, &storeapi).await else { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "stopped", String::new()).await; return Ok(()); }; maintenance_features = features; @@ -2266,16 +2283,19 @@ async fn run_data_scanner_with_maintenance_state( } => (cycle, leader_epoch, revision), ScannerCycleStateStartup::Blocked => { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "stopped", String::new()).await; return Ok(()); } ScannerCycleStateStartup::Transient(err) => { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "stopped", String::new()).await; return Err(err); } }; let (usage_floor, usage_floor_startup) = match persisted_usage_floor_for_startup(storeapi.clone(), true).await { Ok(result) => result, Err(err) => { + let error = err.to_string(); error!( target: "rustfs::scanner", event = EVENT_SCANNER_PERSIST_STATE, @@ -2286,7 +2306,9 @@ async fn run_data_scanner_with_maintenance_state( error = %err, "Scanner stopped because the persisted usage floor could not be loaded" ); + record_scanner_usage_floor_failure(error.clone()); global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "usage_floor_load_failed", error).await; return Ok(()); } }; @@ -2294,10 +2316,13 @@ async fn run_data_scanner_with_maintenance_state( prepare_cycle_for_usage_floor_bootstrap(&mut cycle_info, usage_floor, usage_floor_startup); apply_persisted_usage_floor(&mut cycle_info, &mut leader_epoch, usage_floor); match usage_floor_startup { - PersistedUsageFloorStartup::Authoritative | PersistedUsageFloorStartup::BootstrapPending => {} + PersistedUsageFloorStartup::Authoritative + | PersistedUsageFloorStartup::BootstrapPending + | PersistedUsageFloorStartup::RecoveredLegacyEmptyFence => {} PersistedUsageFloorStartup::Missing => { if ctx.is_cancelled() || guard.is_lock_lost() { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } @@ -2322,10 +2347,12 @@ async fn run_data_scanner_with_maintenance_state( "Scanner stopped because the usage baseline bootstrap could not be initialized" ); global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "usage_floor_bootstrap_failed", err.to_string()).await; return Ok(()); } None => { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } } @@ -2334,6 +2361,7 @@ async fn run_data_scanner_with_maintenance_state( if ctx.is_cancelled() || guard.is_lock_lost() { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } let claim_ctx = ctx.child_token(); @@ -2355,6 +2383,7 @@ async fn run_data_scanner_with_maintenance_state( if guard.is_lock_lost() { record_scanner_leader_lock_lost("Scanner leader lock lost while claiming the leadership epoch").await; global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(true, "lost", String::new()).await; return Ok(()); } if !leadership_claimed { @@ -2367,10 +2396,26 @@ async fn run_data_scanner_with_maintenance_state( state = "epoch_claim_failed", "Scanner stopped because the leadership epoch could not be claimed" ); - global_metrics() - .record_scanner_leader_liveness("epoch_claim_failed", false, "leadership epoch claim failed") - .await; global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "epoch_claim_failed", "leadership epoch claim failed".to_string()).await; + return Ok(()); + } + if usage_floor_startup == PersistedUsageFloorStartup::RecoveredLegacyEmptyFence + && let Err(err) = complete_legacy_empty_usage_floor_recovery(storeapi.clone(), leader_epoch).await + { + let error = err.to_string(); + warn!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + state = "usage_floor_recovery_cleanup_deferred", + path = %DATA_USAGE_RECOVERY_PATH.as_str(), + error = %err, + "Scanner usage floor recovery marker cleanup was deferred" + ); + global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(false, "usage_floor_recovery_pending", error).await; return Ok(()); } @@ -2382,6 +2427,7 @@ async fn run_data_scanner_with_maintenance_state( if guard.is_lock_lost() { record_scanner_leader_lock_lost("Scanner leader lock lost before the initial cycle").await; global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(true, "lost", String::new()).await; return Ok(()); } let cycle_ctx = ctx.child_token(); @@ -2405,10 +2451,12 @@ async fn run_data_scanner_with_maintenance_state( ScannerCycleWaitOutcome::LockLost => { record_scanner_leader_lock_lost("Scanner leader lock lost during the initial cycle").await; global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(true, "lost", String::new()).await; return Ok(()); } ScannerCycleWaitOutcome::Cancelled => { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } ScannerCycleWaitOutcome::Deadline { worker_stopped } => { @@ -2425,6 +2473,7 @@ async fn run_data_scanner_with_maintenance_state( &mut guard, ) .await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } }; @@ -2434,6 +2483,7 @@ async fn run_data_scanner_with_maintenance_state( if guard.is_lock_lost() { record_scanner_leader_lock_lost("Scanner leader lock lost during the initial cycle").await; global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(true, "lost", String::new()).await; return Ok(()); } let runtime_config = resolve_scanner_runtime_config(); @@ -2666,10 +2716,12 @@ async fn run_data_scanner_with_maintenance_state( ScannerCycleWaitOutcome::LockLost => { record_scanner_leader_lock_lost("Scanner leader lock lost during a scanner cycle").await; global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(true, "lost", String::new()).await; return Ok(()); } ScannerCycleWaitOutcome::Cancelled => { global_metrics().set_cycle(None).await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } ScannerCycleWaitOutcome::Deadline { worker_stopped } => { @@ -2686,6 +2738,7 @@ async fn run_data_scanner_with_maintenance_state( &mut guard, ) .await; + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; return Ok(()); } }; @@ -2771,10 +2824,7 @@ async fn run_data_scanner_with_maintenance_state( } global_metrics().set_cycle(None).await; - reset_scanner_cycle_schedule(); - if !guard.is_lock_lost() { - global_metrics().record_scanner_leader_liveness("stopped", false, "").await; - } + finish_scanner_leader_iteration(guard.is_lock_lost(), "stopped", String::new()).await; debug!( target: "rustfs::scanner", diff --git a/crates/scanner/src/scanner/cycle_state.rs b/crates/scanner/src/scanner/cycle_state.rs index 95a42e37f..f34de9317 100644 --- a/crates/scanner/src/scanner/cycle_state.rs +++ b/crates/scanner/src/scanner/cycle_state.rs @@ -14,7 +14,7 @@ /// Scanner cycle-state codec, persisted usage floors, and cycle-state persistence. use super::*; use crate::ScannerGetObjectReader; -use crate::data_usage_define::DATA_USAGE_BLOOM_RECOVERY_PATH; +use crate::data_usage_define::{DATA_USAGE_BLOOM_RECOVERY_PATH, DATA_USAGE_RECOVERY_PATH}; use crate::storage_api::owner::ObjectIO as _; use tokio::io::AsyncReadExt as _; @@ -23,6 +23,8 @@ const MAX_SCANNER_CYCLE_STATE_BYTES: u64 = 1024 * 1024; pub(super) const MAX_SCANNER_CYCLE_RECOVERY_RETRIES: u32 = 5; const METRIC_SCANNER_CYCLE_RECOVERY_REQUIRED: &str = "rustfs_scanner_cycle_recovery_required"; const METRIC_SCANNER_CYCLE_RECOVERY_RETRY_COUNT: &str = "rustfs_scanner_cycle_recovery_retry_count"; +const USAGE_FLOOR_LOAD_FAILED: &str = "usage_floor_load_failed"; +const LEGACY_EMPTY_USAGE_FLOOR_RECOVERY: &str = "legacy_empty_usage_floor"; #[derive(Clone, Debug, Default, Serialize)] pub struct ScannerCycleRecoveryStatus { @@ -62,7 +64,15 @@ pub fn scanner_cycle_recovery_status() -> ScannerCycleRecoveryStatus { } fn set_scanner_cycle_recovery_status(status: ScannerCycleRecoveryStatus) { - let recovery_required = if matches!(status.state.as_str(), "blocked" | "paused" | "recovery-required" | "cleanup-pending") { + let recovery_required = if matches!( + status.state.as_str(), + "blocked" + | "paused" + | "recovery-required" + | "cleanup-pending" + | "usage_floor_load_failed" + | "usage_floor_recovery_pending" + ) { 1.0 } else { 0.0 @@ -74,6 +84,67 @@ fn set_scanner_cycle_recovery_status(status: ScannerCycleRecoveryStatus) { .unwrap_or_else(|poisoned| poisoned.into_inner()) = status; } +pub(super) fn record_scanner_usage_floor_failure(reason: String) { + let previous = scanner_cycle_recovery_status(); + let same_failure = previous.classification.as_deref() == Some(USAGE_FLOOR_LOAD_FAILED); + let now = unix_now_secs(); + let (first_detected_at_unix_secs, retry_count) = if same_failure { + (previous.first_detected_at_unix_secs.or(Some(now)), previous.retry_count) + } else { + (Some(now), 0) + }; + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_OBJ_NAME_PATH.clone(), + state: USAGE_FLOOR_LOAD_FAILED.to_string(), + classification: Some(USAGE_FLOOR_LOAD_FAILED.to_string()), + first_detected_at_unix_secs, + last_attempt_at_unix_secs: Some(now), + retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: true, + reason: Some(reason), + ..Default::default() + }); +} + +pub(super) fn clear_scanner_usage_floor_failure() { + if scanner_cycle_recovery_status().classification.as_deref() == Some(USAGE_FLOOR_LOAD_FAILED) { + set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); + } +} + +pub(super) fn record_legacy_empty_usage_floor_recovery_pending(leader_epoch: u64) { + let previous = scanner_cycle_recovery_status(); + let same_recovery = previous.classification.as_deref() == Some(LEGACY_EMPTY_USAGE_FLOOR_RECOVERY) + && previous.leader_epoch == Some(leader_epoch); + let now = unix_now_secs(); + let (first_detected_at_unix_secs, retry_count) = if same_recovery { + (previous.first_detected_at_unix_secs.or(Some(now)), previous.retry_count) + } else { + (Some(now), 0) + }; + set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus { + path: DATA_USAGE_OBJ_NAME_PATH.clone(), + quarantine_path: Some(DATA_USAGE_RECOVERY_PATH.clone()), + state: "usage_floor_recovery_pending".to_string(), + classification: Some(LEGACY_EMPTY_USAGE_FLOOR_RECOVERY.to_string()), + leader_epoch: Some(leader_epoch), + first_detected_at_unix_secs, + last_attempt_at_unix_secs: Some(now), + retry_count, + max_retries: MAX_SCANNER_CYCLE_RECOVERY_RETRIES, + retryable: true, + reason: Some("legacy empty usage floor recovery is awaiting a fenced leadership claim".to_string()), + ..Default::default() + }); +} + +pub(super) fn clear_legacy_empty_usage_floor_recovery_status() { + if scanner_cycle_recovery_status().classification.as_deref() == Some(LEGACY_EMPTY_USAGE_FLOOR_RECOVERY) { + set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); + } +} + pub(super) fn record_scanner_cycle_recovery_retry(attempt: u32) -> bool { let mut status = scanner_cycle_recovery_status(); status.retry_count = u64::from(attempt); @@ -1185,6 +1256,246 @@ pub(super) enum PersistedUsageFloorStartup { Authoritative, Missing, BootstrapPending, + RecoveredLegacyEmptyFence, +} + +#[derive(Clone, Debug)] +struct LegacyEmptyUsageFloorPrimary { + revision: DataUsageCacheRevision, + epoch: u64, +} + +#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] +#[serde(deny_unknown_fields)] +struct LegacyEmptyUsageFloorRecoveryMarker { + schema_version: u16, + primary_revision: String, + leader_epoch: u64, +} + +async fn read_legacy_empty_usage_floor_recovery_marker( + storeapi: Arc, +) -> Result, ScannerError> { + let (data, revision) = read_config_with_revision(storeapi, DATA_USAGE_RECOVERY_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to read scanner usage recovery marker: {err}")))?; + let Some(data) = data else { + return Ok(None); + }; + let marker = serde_json::from_slice::(&data) + .map_err(|err| ScannerError::Other(format!("failed to decode scanner usage recovery marker: {err}")))?; + if marker.schema_version != 1 || marker.primary_revision.is_empty() || marker.leader_epoch == 0 { + return Err(ScannerError::Other("scanner usage recovery marker is invalid".to_string())); + } + if !matches!(revision, DataUsageCacheRevision::Etag(_)) { + return Err(ScannerError::Other("scanner usage recovery marker has no revision".to_string())); + } + Ok(Some((marker, revision))) +} + +async fn clear_legacy_empty_usage_floor_recovery_marker( + storeapi: Arc, + marker_revision: &DataUsageCacheRevision, + expected_publication_epoch: u64, +) -> Result<(), ScannerError> { + let delete_result = delete_config_with_publication_admission_for_epoch( + storeapi.clone(), + RUSTFS_META_BUCKET, + DATA_USAGE_RECOVERY_PATH.as_str(), + ScannerObjectOptions { + delete_prefix: false, + http_preconditions: Some(marker_revision.preconditions()), + ..Default::default() + }, + expected_publication_epoch, + ) + .await; + match delete_result { + Ok(_) => Ok(()), + Err(err) => { + let (_, revision) = read_config_with_revision(storeapi, DATA_USAGE_RECOVERY_PATH.as_str()) + .await + .map_err(|read_err| { + ScannerError::Other(format!("failed to reconcile scanner usage recovery cleanup: {read_err}")) + })?; + if matches!(revision, DataUsageCacheRevision::Missing) { + Ok(()) + } else { + Err(ScannerError::Other(format!("failed to clear scanner usage recovery marker: {err}"))) + } + } + } +} + +pub(super) async fn complete_legacy_empty_usage_floor_recovery( + storeapi: Arc, + claimed_epoch: u64, +) -> Result<(), ScannerError> { + let Some((marker, marker_revision)) = read_legacy_empty_usage_floor_recovery_marker(storeapi.clone()).await? else { + return Ok(()); + }; + if claimed_epoch <= marker.leader_epoch { + return Err(ScannerError::Other("scanner usage recovery did not advance the leader epoch".to_string())); + } + let (primary, _) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to verify recovered scanner usage bootstrap: {err}")))?; + let primary = primary.ok_or_else(|| ScannerError::Other("recovered scanner usage bootstrap is missing".to_string()))?; + let usage = serde_json::from_slice::(&primary) + .map_err(|err| ScannerError::Other(format!("failed to decode recovered scanner usage bootstrap: {err}")))?; + if !data_usage_info_is_bootstrap_pending(&usage) || usage.scanner_epoch != Some(claimed_epoch) { + return Err(ScannerError::Other( + "recovered scanner usage bootstrap does not match the claimed epoch".to_string(), + )); + } + let expected_publication_epoch = scanner_publication_epoch(storeapi.clone()) + .await + .ok_or_else(|| ScannerError::Other("scanner usage recovery cleanup is blocked by data movement".to_string()))?; + clear_legacy_empty_usage_floor_recovery_marker(storeapi, &marker_revision, expected_publication_epoch).await?; + clear_legacy_empty_usage_floor_recovery_status(); + Ok(()) +} + +fn legacy_empty_usage_fence_epoch(data: &[u8], usage: &DataUsageInfo) -> Option> { + if usage.last_update.is_none() || usage.scanner_cycle.is_some() { + return None; + } + if usage.scanner_epoch.is_some_and(|epoch| epoch == 0 || epoch >= u64::MAX - 1) { + return None; + } + let expected = DataUsageInfo { + last_update: usage.last_update, + scanner_epoch: usage.scanner_epoch, + ..Default::default() + }; + if usage != &expected { + return None; + } + + let serde_json::Value::Object(fields) = serde_json::from_slice::(data).ok()? else { + return None; + }; + // RUSTFS_COMPAT_TODO(backlog-2102): accept only the exact empty usage fence serialized by rc.2/rc.3. Remove after those releases are no longer supported direct-upgrade sources. + const REQUIRED_FIELDS: &[&str] = &[ + "total_capacity", + "total_used_capacity", + "total_free_capacity", + "last_update", + "objects_total_count", + "versions_total_count", + "delete_markers_total_count", + "objects_total_size", + "replication_info", + "buckets_count", + "buckets_usage", + "usage_snapshot_complete", + "bucket_sizes", + "disk_usage_status", + ]; + let expected_len = REQUIRED_FIELDS.len() + if usage.scanner_epoch.is_some() { 1 } else { 0 }; + if fields.len() != expected_len + || REQUIRED_FIELDS.iter().any(|field| !fields.contains_key(*field)) + || (usage.scanner_epoch.is_some() != fields.contains_key("scanner_epoch")) + { + return None; + } + Some(usage.scanner_epoch) +} + +async fn recover_legacy_empty_usage_floor( + storeapi: Arc, + primary: LegacyEmptyUsageFloorPrimary, + expected_publication_epoch: u64, +) -> Result<(), ScannerError> { + let DataUsageCacheRevision::Etag(primary_revision) = &primary.revision else { + return Err(ScannerError::Other("legacy empty scanner usage floor has no revision".to_string())); + }; + let marker = LegacyEmptyUsageFloorRecoveryMarker { + schema_version: 1, + primary_revision: primary_revision.clone(), + leader_epoch: primary.epoch, + }; + let marker_data = serde_json::to_vec(&marker) + .map_err(|err| ScannerError::Other(format!("failed to encode scanner usage recovery marker: {err}")))?; + match read_legacy_empty_usage_floor_recovery_marker(storeapi.clone()).await? { + Some((persisted, _)) if persisted != marker => { + return Err(ScannerError::Other( + "scanner usage recovery marker conflicts with the persisted empty floor".to_string(), + )); + } + Some(_) => {} + None => { + let marker_save = save_config_with_publication_admission_for_epoch( + storeapi.clone(), + DATA_USAGE_RECOVERY_PATH.as_str(), + marker_data.clone(), + DataUsageCacheRevision::Missing.preconditions(), + expected_publication_epoch, + ) + .await; + if !marker_save + .as_ref() + .ok() + .and_then(|info| info.etag.as_deref()) + .is_some_and(|etag| !etag.is_empty()) + { + let persisted = read_legacy_empty_usage_floor_recovery_marker(storeapi.clone()).await?; + if persisted.as_ref().map(|(persisted, _)| persisted) != Some(&marker) { + return Err(ScannerError::Other(match marker_save { + Ok(_) => "scanner usage recovery marker returned no ETag and could not be confirmed".to_string(), + Err(err) => format!("failed to persist scanner usage recovery marker: {err}"), + })); + } + } + } + } + + let marker = DataUsageInfo { + last_update: Some(std::time::SystemTime::now()), + scanner_epoch: Some(primary.epoch), + usage_snapshot_converged: Some(false), + usage_snapshot_bootstrap_pending: true, + ..Default::default() + }; + let data = serde_json::to_vec(&marker) + .map_err(|err| ScannerError::Other(format!("failed to encode recovered scanner usage bootstrap: {err}")))?; + let save_result = save_config_with_publication_admission_for_epoch( + storeapi.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + data.clone(), + primary.revision.preconditions(), + expected_publication_epoch, + ) + .await; + if save_result + .as_ref() + .ok() + .and_then(|info| info.etag.as_deref()) + .is_some_and(|etag| !etag.is_empty()) + { + warn!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + state = "legacy_empty_usage_floor_recovered", + path = %DATA_USAGE_OBJ_NAME_PATH.as_str(), + scanner_epoch = primary.epoch, + "Scanner recovered a legacy empty usage floor" + ); + return Ok(()); + } + + let (persisted, revision) = read_config_with_revision(storeapi, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .map_err(|err| ScannerError::Other(format!("failed to reconcile recovered 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(_) => "recovered scanner usage bootstrap returned no ETag and could not be confirmed".to_string(), + Err(err) => format!("failed to recover legacy empty scanner usage floor: {err}"), + })) } pub(super) fn encode_scanner_cycle_state( @@ -1324,16 +1635,22 @@ pub(super) async fn persisted_usage_floor_for_startup( 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 recovery_marker = read_legacy_empty_usage_floor_recovery_marker(storeapi.clone()).await?; let mut floor = PersistedUsageFloor::default(); let mut found_any = false; let mut bootstrap_pending = false; + let mut recovered_bootstrap = false; + let mut bootstrap_epoch = None; // A valid JSON object without a baseline identity is not a floor and must // never be treated as an empty one. It can, however, be a partially // written v2 primary left behind during an upgrade. Keep its epoch as a // fence while looking for a durable companion snapshot; if no companion // is new enough, the caller still fails closed below. let mut invalid_baseline_path: Option = None; - let mut invalid_baseline_epoch: Option = None; + let mut invalid_baseline_epoch = recovery_marker.as_ref().map(|(marker, _)| marker.leader_epoch); + let mut unrecoverable_baseline_path: Option = None; + let mut stale_authoritative_path: Option = None; + let mut legacy_empty_primary: Option = None; 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 { @@ -1348,8 +1665,9 @@ pub(super) async fn persisted_usage_floor_for_startup( for primary_path in [DATA_USAGE_OBJ_NAME_PATH.as_str(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()] { let backup_path = format!("{primary_path}.bkp"); let is_v2_path = primary_path == DATA_USAGE_OBJ_NAME_PATH.as_str(); + let mut recovered_primary_companion_epoch = None; let primary_epoch = match read_config_with_revision(storeapi.clone(), primary_path).await { - Ok((Some(data), _)) => { + Ok((Some(data), revision)) => { let usage = serde_json::from_slice::(&data).map_err(|err| { ScannerError::Other(format!("failed to decode scanner usage floor from {primary_path}: {err}")) })?; @@ -1358,18 +1676,47 @@ pub(super) async fn persisted_usage_floor_for_startup( return Err(ScannerError::Other("multiple scanner usage bootstrap markers were found".to_string())); } bootstrap_pending = true; + bootstrap_epoch = usage.scanner_epoch; + if let Some((marker, _)) = recovery_marker.as_ref() { + if usage.scanner_epoch.is_none_or(|epoch| epoch < marker.leader_epoch) { + return Err(ScannerError::Other( + "scanner usage bootstrap is older than its recovery marker".to_string(), + )); + } + recovered_bootstrap = true; + } update_floor(&mut floor, &usage, primary_path)?; None } else if !data_usage_info_has_persisted_baseline_identity(&usage) { invalid_baseline_path.get_or_insert_with(|| primary_path.to_string()); invalid_baseline_epoch = invalid_baseline_epoch.max(usage.scanner_epoch); + match legacy_empty_usage_fence_epoch(&data, &usage) { + Some(Some(epoch)) if is_v2_path => { + legacy_empty_primary = Some(LegacyEmptyUsageFloorPrimary { revision, epoch }); + } + Some(_) => {} + None => { + unrecoverable_baseline_path.get_or_insert_with(|| primary_path.to_string()); + } + } None } else { let epoch = usage.scanner_epoch.unwrap_or_default(); + if recovered_bootstrap && !is_v2_path { + if epoch < floor.leader_epoch { + unrecoverable_baseline_path.get_or_insert_with(|| primary_path.to_string()); + } else { + update_floor(&mut floor, &usage, primary_path)?; + recovered_primary_companion_epoch = Some(epoch); + } + None // A legacy snapshot may be structurally valid but older // than an incomplete v2 snapshot left by a newer leader. // Do not let that candidate regress the startup floor. - if !is_v2_path && invalid_baseline_epoch.is_some_and(|fenced_epoch| epoch < fenced_epoch) { + } else if invalid_baseline_epoch.is_some_and(|fenced_epoch| epoch < fenced_epoch) + && (!is_v2_path || recovery_marker.is_some()) + { + stale_authoritative_path.get_or_insert_with(|| primary_path.to_string()); None } else { update_floor(&mut floor, &usage, primary_path)?; @@ -1387,17 +1734,51 @@ pub(super) async fn persisted_usage_floor_for_startup( let mut any_found = primary_epoch.is_some(); match read_config_with_revision(storeapi.clone(), &backup_path).await { Ok((Some(data), _)) => { + let usage = serde_json::from_slice::(&data).map_err(|err| { + ScannerError::Other(format!("failed to decode scanner usage floor from {backup_path}: {err}")) + })?; if bootstrap_pending { + if let Some(primary_epoch) = recovered_primary_companion_epoch { + if data_usage_info_has_persisted_baseline_identity(&usage) { + let backup_epoch = usage.scanner_epoch.unwrap_or_default(); + if backup_epoch >= primary_epoch { + update_floor(&mut floor, &usage, &backup_path)?; + } + } else if let Some(epoch) = legacy_empty_usage_fence_epoch(&data, &usage) { + if let Some(epoch) = epoch { + floor.leader_epoch = floor.leader_epoch.max(epoch); + } + } else { + unrecoverable_baseline_path.get_or_insert_with(|| backup_path.clone()); + } + continue; + } + let compatible_empty_fence = legacy_empty_usage_fence_epoch(&data, &usage).is_some_and(|epoch| { + epoch.is_none_or(|epoch| bootstrap_epoch.is_some_and(|bootstrap_epoch| epoch <= bootstrap_epoch)) + }); + if compatible_empty_fence { + continue; + } + if recovered_bootstrap && data_usage_info_has_persisted_baseline_identity(&usage) { + let epoch = usage.scanner_epoch.unwrap_or_default(); + if epoch >= floor.leader_epoch { + update_floor(&mut floor, &usage, &backup_path)?; + if primary_path == DATA_USAGE_OBJ_NAME_PATH.as_str() { + break; + } + continue; + } + } return Err(ScannerError::Other( "scanner usage bootstrap conflicts with a persisted backup".to_string(), )); } - let usage = serde_json::from_slice::(&data).map_err(|err| { - ScannerError::Other(format!("failed to decode scanner usage floor from {backup_path}: {err}")) - })?; if !data_usage_info_has_persisted_baseline_identity(&usage) { invalid_baseline_path.get_or_insert_with(|| backup_path.clone()); invalid_baseline_epoch = invalid_baseline_epoch.max(usage.scanner_epoch); + if legacy_empty_usage_fence_epoch(&data, &usage).is_none() { + unrecoverable_baseline_path.get_or_insert_with(|| backup_path.clone()); + } // This is still persisted state, so it must not enable a // missing-state bootstrap. Continue to a legacy pair in // case it contains a complete, fenced snapshot. @@ -1411,6 +1792,8 @@ pub(super) async fn persisted_usage_floor_for_startup( { update_floor(&mut floor, &usage, &backup_path)?; any_found = true; + } else { + stale_authoritative_path.get_or_insert_with(|| backup_path.clone()); } } } @@ -1433,6 +1816,30 @@ pub(super) async fn persisted_usage_floor_for_startup( } if !found_any && !bootstrap_pending { + if allow_missing_for_bootstrap + && unrecoverable_baseline_path.is_none() + && stale_authoritative_path.is_none() + && let Some(mut primary) = legacy_empty_primary + { + primary.epoch = primary.epoch.max(invalid_baseline_epoch.unwrap_or_default()); + recover_legacy_empty_usage_floor(storeapi.clone(), primary.clone(), read_epoch).await?; + record_legacy_empty_usage_floor_recovery_pending(primary.epoch); + return Ok(( + PersistedUsageFloor { + next_cycle: 0, + leader_epoch: primary.epoch, + }, + PersistedUsageFloorStartup::RecoveredLegacyEmptyFence, + )); + } + if let Some(path) = stale_authoritative_path { + return Err(ScannerError::Other(format!( + "persisted scanner usage floor from {path} is older than the required recovery fence" + ))); + } + if recovery_marker.is_some() { + return Err(ScannerError::Other("scanner usage recovery marker has no matching primary".to_string())); + } if let Some(path) = invalid_baseline_path { return Err(ScannerError::Other(format!( "persisted scanner usage floor from {path} has no authoritative baseline or newer valid backup" @@ -1470,18 +1877,67 @@ pub(super) async fn persisted_usage_floor_for_startup( } drop(publication_admission); } - let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi, read_epoch).await else { + let Some(publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { return Err(ScannerError::Other( "scanner usage floor changed while its epoch proof was being confirmed".to_string(), )); }; let state = if found_any { PersistedUsageFloorStartup::Authoritative + } else if recovered_bootstrap { + if let Some(path) = unrecoverable_baseline_path { + return Err(ScannerError::Other(format!( + "scanner usage recovery conflicts with persisted usage state at {path}" + ))); + } + let recovery_epoch = recovery_marker + .as_ref() + .map(|(marker, _)| marker.leader_epoch) + .unwrap_or(floor.leader_epoch); + record_legacy_empty_usage_floor_recovery_pending(recovery_epoch); + PersistedUsageFloorStartup::RecoveredLegacyEmptyFence } else if bootstrap_pending { + if let Some(path) = unrecoverable_baseline_path { + return Err(ScannerError::Other(format!( + "scanner usage bootstrap conflicts with persisted usage state at {path}" + ))); + } PersistedUsageFloorStartup::BootstrapPending } else { PersistedUsageFloorStartup::Missing }; + if found_any && let Some((_, marker_revision)) = recovery_marker.as_ref() { + drop(publication_admission); + let marker_cleared = + match clear_legacy_empty_usage_floor_recovery_marker(storeapi.clone(), marker_revision, read_epoch).await { + Ok(()) => true, + Err(err) => { + warn!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + state = "usage_floor_recovery_cleanup_deferred", + path = %DATA_USAGE_RECOVERY_PATH.as_str(), + error = %err, + "Scanner usage floor recovery marker cleanup was deferred" + ); + false + } + }; + let Some(_final_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await + else { + return Err(ScannerError::Other( + "scanner usage floor changed after recovery marker cleanup".to_string(), + )); + }; + if marker_cleared { + clear_legacy_empty_usage_floor_recovery_status(); + } + clear_scanner_usage_floor_failure(); + return Ok((floor, state)); + } + clear_scanner_usage_floor_failure(); Ok((floor, state)) } diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 5cc87c063..b0ea9c202 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -361,6 +361,7 @@ struct MemoryConfigStore { cancel_after_interleaving_puts: Mutex>, cancel_after_successful_puts: Mutex>, replace_after_successful_puts: Mutex)>>, + error_after_commit_deletes: Mutex>, put_counts: Mutex>, publication_admission_blocked: AtomicBool, block_publication_after_admissions: AtomicUsize, @@ -1959,6 +1960,562 @@ async fn scanner_usage_floor_keeps_valid_primary_when_backup_has_no_identity() { ); } +fn rc3_legacy_empty_usage_fence(epoch: Option) -> Vec { + // Pinned field set emitted by rc.3 after DeleteBucket synthesized a + // default v2 usage primary. Leadership added scanner_epoch separately. + const RC3_EMPTY_USAGE_FENCE: &str = r#"{ + "total_capacity":0, + "total_used_capacity":0, + "total_free_capacity":0, + "last_update":{"secs_since_epoch":1,"nanos_since_epoch":0}, + "objects_total_count":0, + "versions_total_count":0, + "delete_markers_total_count":0, + "objects_total_size":0, + "replication_info":{}, + "buckets_count":0, + "buckets_usage":{}, + "usage_snapshot_complete":false, + "bucket_sizes":{}, + "disk_usage_status":[] + }"#; + let mut value = + serde_json::from_str::(RC3_EMPTY_USAGE_FENCE).expect("pinned rc.3 empty usage fence should decode"); + let fields = value + .as_object_mut() + .expect("legacy empty usage fence should be a JSON object"); + if let Some(epoch) = epoch { + fields.insert("scanner_epoch".to_string(), serde_json::Value::from(epoch)); + } + serde_json::to_vec(&value).expect("rc.3 legacy empty usage fence fixture should encode") +} + +#[tokio::test] +async fn scanner_usage_floor_recovers_rc3_empty_fences_and_restarts_from_zero() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); + let backup_key = memory_config_key(RUSTFS_META_BUCKET, &backup_path); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store + .objects + .lock() + .await + .insert(backup_key, rc3_legacy_empty_usage_fence(None)); + store.revisions.lock().await.insert(primary_key, 1); + + let (floor, startup) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("rc.3 empty usage fences should enter recovery"); + assert_eq!( + floor, + PersistedUsageFloor { + next_cycle: 0, + leader_epoch: 7, + } + ); + assert_eq!(startup, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + + let primary = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("recovered usage bootstrap should be persisted"); + let pending = serde_json::from_slice::(&primary).expect("recovered usage bootstrap should decode"); + assert!(data_usage_info_is_bootstrap_pending(&pending)); + assert_eq!(pending.scanner_epoch, Some(7)); + assert!(read_config(store.clone(), DATA_USAGE_RECOVERY_PATH.as_str()).await.is_ok()); + + let (restart_floor, restart_state) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("recovery marker should survive a restart before leadership claim"); + assert_eq!(restart_floor.leader_epoch, 7); + assert_eq!(restart_state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + let mut cycle = CurrentCycle { + current: 41, + next: 42, + ..Default::default() + }; + assert_eq!( + prepare_cycle_for_usage_floor_bootstrap(&mut cycle, restart_floor, restart_state), + (true, true) + ); + assert_eq!(cycle.current, 0); + assert_eq!(cycle.next, 0); + assert!(cycle.cycle_completed.is_empty()); + + let mut revision = DataUsageCacheRevision::Missing; + let mut leader_epoch = restart_floor.leader_epoch; + assert!( + claim_scanner_leadership( + &CancellationToken::new(), + store.clone(), + &mut cycle, + &mut revision, + &mut leader_epoch, + true, + true, + ) + .await + ); + assert_eq!(leader_epoch, 8); + complete_legacy_empty_usage_floor_recovery(store.clone(), leader_epoch) + .await + .expect("leadership claim should retire the recovery marker"); + assert!(matches!( + read_config(store.clone(), DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); + let (claimed_floor, claimed_state) = persisted_usage_floor_for_startup(store, true) + .await + .expect("claimed bootstrap should remain restartable"); + assert_eq!(claimed_floor.leader_epoch, 8); + assert_eq!(claimed_state, PersistedUsageFloorStartup::BootstrapPending); +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_preserves_newer_authoritative_companion_floor() { + for companion_path in [ + format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()), + LEGACY_DATA_USAGE_OBJ_NAME_PATH.clone(), + ] { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key, 1); + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("legacy empty primary should enter recovery"); + + let mut companion = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + companion.scanner_epoch = Some(8); + companion.scanner_cycle = Some(11); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, &companion_path), + serde_json::to_vec(&companion).expect("authoritative companion should encode"), + ); + if companion_path.ends_with(".bkp") { + let mut stale_legacy = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + stale_legacy.scanner_epoch = Some(6); + stale_legacy.scanner_cycle = Some(10); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()), + serde_json::to_vec(&stale_legacy).expect("stale legacy companion should encode"), + ); + } else { + let mut stale_backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + stale_backup.scanner_epoch = Some(6); + stale_backup.scanner_cycle = Some(10); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, &format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str())), + serde_json::to_vec(&stale_backup).expect("stale legacy backup should encode"), + ); + } + + let (floor, state) = persisted_usage_floor_for_startup(store, true) + .await + .expect("a newer authoritative companion should advance the recovery floor"); + assert_eq!(floor.leader_epoch, 8, "unexpected companion path: {companion_path}"); + assert_eq!(floor.next_cycle, 12, "unexpected companion path: {companion_path}"); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + } +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_fences_non_authoritative_legacy_backup() { + for partial_backup in [false, true] { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key, 1); + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("legacy empty primary should enter recovery"); + + let mut legacy_primary = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + legacy_primary.scanner_epoch = Some(8); + legacy_primary.scanner_cycle = Some(11); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()), + serde_json::to_vec(&legacy_primary).expect("legacy primary should encode"), + ); + let backup_path = format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()); + let backup = if partial_backup { + serde_json::to_vec(&DataUsageInfo { + last_update: Some(std::time::SystemTime::UNIX_EPOCH), + scanner_epoch: Some(9), + buckets_count: 1, + ..Default::default() + }) + .expect("partial legacy backup should encode") + } else { + rc3_legacy_empty_usage_fence(Some(9)) + }; + store + .objects + .lock() + .await + .insert(memory_config_key(RUSTFS_META_BUCKET, &backup_path), backup); + + if partial_backup { + let err = persisted_usage_floor_for_startup(store, true) + .await + .expect_err("a partial noncanonical backup must remain fail-closed"); + assert!(err.to_string().contains("conflicts with persisted usage state")); + } else { + let (floor, state) = persisted_usage_floor_for_startup(store, true) + .await + .expect("an exact empty backup should contribute its epoch fence"); + assert_eq!(floor.leader_epoch, 9); + assert_eq!(floor.next_cycle, 12); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + } + } +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_resumes_after_marker_only_crash_point() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + let original = rc3_legacy_empty_usage_fence(Some(7)); + store.objects.lock().await.insert(primary_key.clone(), original.clone()); + store.revisions.lock().await.insert(primary_key.clone(), 1); + store.fail_put_number.lock().await.insert(primary_key, 1); + + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("injected primary CAS failure should leave recovery pending"); + assert!(read_config(store.clone(), DATA_USAGE_RECOVERY_PATH.as_str()).await.is_ok()); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("legacy primary should remain after the failed CAS"), + original + ); + + let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("the durable marker should resume the primary conversion"); + assert_eq!(floor.leader_epoch, 7); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + let recovered = read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("recovered bootstrap should replace the legacy primary"); + assert!(data_usage_info_is_bootstrap_pending( + &serde_json::from_slice(&recovered).expect("recovered bootstrap should decode") + )); +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_reconciles_marker_post_commit_error() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + let marker_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_RECOVERY_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key, 1); + store.error_after_commit_put_number.lock().await.insert(marker_key, 1); + + let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("a committed recovery marker should reconcile after an ambiguous error"); + assert_eq!(floor.leader_epoch, 7); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + assert!(read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await.is_ok()); +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_reconciles_marker_delete_post_commit_error() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key, 1); + let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("legacy empty primary should enter recovery"); + let mut cycle = CurrentCycle::default(); + let mut revision = DataUsageCacheRevision::Missing; + let mut leader_epoch = floor.leader_epoch; + let (allow_pending, reset_on_conflict) = prepare_cycle_for_usage_floor_bootstrap(&mut cycle, floor, state); + assert!( + claim_scanner_leadership( + &CancellationToken::new(), + store.clone(), + &mut cycle, + &mut revision, + &mut leader_epoch, + allow_pending, + reset_on_conflict, + ) + .await + ); + store + .error_after_commit_deletes + .lock() + .await + .insert(memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_RECOVERY_PATH.as_str())); + + complete_legacy_empty_usage_floor_recovery(store.clone(), leader_epoch) + .await + .expect("a committed marker delete should reconcile after an ambiguous error"); + assert!(matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); + assert_eq!(scanner_cycle_recovery_status().state, "healthy"); +} + +#[tokio::test] +#[serial] +async fn scanner_usage_floor_recovery_retry_budget_uses_marker_epoch_identity() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key.clone(), 1); + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("legacy empty primary should enter recovery"); + assert!(record_scanner_cycle_recovery_retry(3)); + let first_detected = scanner_cycle_recovery_status().first_detected_at_unix_secs; + + let primary = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("recovered bootstrap should exist"); + let mut pending = serde_json::from_slice::(&primary).expect("recovered bootstrap should decode"); + pending.scanner_epoch = Some(8); + store.objects.lock().await.insert( + primary_key.clone(), + serde_json::to_vec(&pending).expect("claimed bootstrap should encode"), + ); + *store.revisions.lock().await.entry(primary_key).or_insert(1) += 1; + + let (floor, state) = persisted_usage_floor_for_startup(store, true) + .await + .expect("claimed bootstrap should retain its recovery identity"); + assert_eq!(floor.leader_epoch, 8); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + let status = scanner_cycle_recovery_status(); + assert_eq!(status.leader_epoch, Some(7)); + assert_eq!(status.retry_count, 3); + assert_eq!(status.first_detected_at_unix_secs, first_detected); + clear_legacy_empty_usage_floor_recovery_status(); +} + +#[tokio::test] +async fn scanner_usage_floor_rejects_noncanonical_empty_fence() { + let store = Arc::new(MemoryConfigStore::default()); + let mut value = serde_json::from_slice::(&rc3_legacy_empty_usage_fence(Some(7))) + .expect("legacy fixture should decode"); + value + .as_object_mut() + .expect("legacy fixture should be an object") + .insert("future_field".to_string(), serde_json::Value::Bool(true)); + store.objects.lock().await.insert( + memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), + serde_json::to_vec(&value).expect("noncanonical fixture should encode"), + ); + + let err = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("unknown legacy fields must not be recovered as an empty baseline"); + assert!(err.to_string().contains("no authoritative baseline")); + assert!(matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn scanner_usage_floor_rejects_zero_or_exhausted_empty_fence_epoch() { + for epoch in [0, u64::MAX - 1, u64::MAX] { + let store = Arc::new(MemoryConfigStore::default()); + let primary = rc3_legacy_empty_usage_fence(Some(epoch)); + store + .objects + .lock() + .await + .insert(memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()), primary.clone()); + + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("an unclaimable legacy epoch must remain fail-closed"); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("rejected legacy floor should remain unchanged"), + primary + ); + assert!(matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); + } +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_does_not_overwrite_concurrent_authoritative_snapshot() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key.clone(), 1); + let mut authoritative = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + authoritative.scanner_epoch = Some(8); + authoritative.scanner_cycle = Some(11); + let authoritative = serde_json::to_vec(&authoritative).expect("authoritative usage should encode"); + store + .interleaving_puts + .lock() + .await + .insert(primary_key, (1, authoritative.clone())); + + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("recovery CAS must lose to a concurrent authoritative snapshot"); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("concurrent authoritative usage should remain"), + authoritative + ); + + let (floor, state) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("the concurrent authoritative snapshot should win on retry"); + assert_eq!(floor.leader_epoch, 8); + assert_eq!(floor.next_cycle, 12); + assert_eq!(state, PersistedUsageFloorStartup::Authoritative); + assert!(matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn scanner_usage_floor_recovery_rejects_concurrent_authoritative_epoch_regression() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + store + .objects + .lock() + .await + .insert(primary_key.clone(), rc3_legacy_empty_usage_fence(Some(7))); + store.revisions.lock().await.insert(primary_key.clone(), 1); + let mut stale = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0); + stale.scanner_epoch = Some(6); + stale.scanner_cycle = Some(11); + let stale = serde_json::to_vec(&stale).expect("stale authoritative usage should encode"); + store.interleaving_puts.lock().await.insert(primary_key, (1, stale.clone())); + + persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("recovery CAS must lose to the concurrent writer"); + let retry_error = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("the recovery marker must fence an older authoritative winner"); + assert!(retry_error.to_string().contains("older than the required recovery fence")); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("stale concurrent snapshot should not be rewritten without a new scan"), + stale + ); + assert!(read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await.is_ok()); +} + +#[test] +#[serial] +fn scanner_usage_floor_failure_is_exposed_and_cleared() { + record_scanner_usage_floor_failure("persisted usage floor is invalid".to_string()); + let blocked = scanner_cycle_recovery_status(); + assert_eq!(blocked.path, DATA_USAGE_OBJ_NAME_PATH.as_str()); + assert_eq!(blocked.state, "usage_floor_load_failed"); + assert_eq!(blocked.classification.as_deref(), Some("usage_floor_load_failed")); + assert!(blocked.retryable); + assert_eq!(blocked.reason.as_deref(), Some("persisted usage floor is invalid")); + let first_detected = blocked.first_detected_at_unix_secs; + + assert!(record_scanner_cycle_recovery_retry(2)); + record_scanner_usage_floor_failure("persisted usage floor remains invalid".to_string()); + let retried = scanner_cycle_recovery_status(); + assert_eq!(retried.retry_count, 2); + assert_eq!(retried.first_detected_at_unix_secs, first_detected); + + clear_scanner_usage_floor_failure(); + let healthy = scanner_cycle_recovery_status(); + assert_eq!(healthy.state, "healthy"); + assert_eq!(healthy.path, DATA_USAGE_BLOOM_NAME_PATH.as_str()); + assert_eq!(healthy.classification, None); +} + +#[test] +#[serial] +fn scanner_usage_floor_recovery_stays_retryable_until_claim_cleanup() { + record_legacy_empty_usage_floor_recovery_pending(7); + let pending = scanner_cycle_recovery_status(); + assert_eq!(pending.state, "usage_floor_recovery_pending"); + assert_eq!(pending.classification.as_deref(), Some("legacy_empty_usage_floor")); + assert_eq!(pending.leader_epoch, Some(7)); + assert!(pending.retryable); + assert_eq!(pending.quarantine_path.as_deref(), Some(DATA_USAGE_RECOVERY_PATH.as_str())); + let first_detected = pending.first_detected_at_unix_secs; + + assert!(record_scanner_cycle_recovery_retry(3)); + record_legacy_empty_usage_floor_recovery_pending(7); + let retried = scanner_cycle_recovery_status(); + assert_eq!(retried.retry_count, 3); + assert_eq!(retried.first_detected_at_unix_secs, first_detected); + + clear_legacy_empty_usage_floor_recovery_status(); + assert_eq!(scanner_cycle_recovery_status().state, "healthy"); +} + +#[tokio::test] +#[serial] +async fn scanner_usage_floor_failure_clears_stale_leader_liveness() { + record_scanner_cycle_schedule_role("leader"); + global_metrics().record_scanner_leader_liveness("acquired", true, "").await; + + finish_scanner_leader_iteration(false, "usage_floor_load_failed", "invalid floor".to_string()).await; + + assert_eq!(scanner_cycle_schedule_status().execution_role, "unknown"); + let report = global_metrics().report().await; + assert_eq!(report.leader_lock_state, "usage_floor_load_failed"); + assert!(!report.leader_lock_held_by_this_process); + assert_eq!(report.leader_lock_last_error, "invalid floor"); + + global_metrics().record_scanner_leader_liveness("acquired", true, "").await; + finish_scanner_leader_iteration(true, "stopped", "lock lost before classification".to_string()).await; + let report = global_metrics().report().await; + assert_eq!(report.leader_lock_state, "stopped"); + assert!(!report.leader_lock_held_by_this_process); + assert_eq!(report.leader_lock_last_error, "lock lost before classification"); +} + #[tokio::test] async fn scanner_usage_floor_recovers_from_incomplete_v2_primary_using_fenced_backup() { let store = Arc::new(MemoryConfigStore::default()); @@ -2039,7 +2596,7 @@ async fn scanner_usage_floor_rejects_backup_older_than_incomplete_v2_primary() { let err = persisted_usage_floor_for_startup(store, true) .await .expect_err("an older backup must not cross the incomplete primary epoch fence"); - assert!(err.to_string().contains("no authoritative baseline")); + assert!(err.to_string().contains("older than the required recovery fence")); } #[tokio::test] @@ -2067,7 +2624,7 @@ async fn scanner_usage_floor_rejects_older_legacy_primary_after_incomplete_v2_pr let err = persisted_usage_floor_for_startup(store, true) .await .expect_err("an older legacy baseline must not cross the incomplete v2 epoch fence"); - assert!(err.to_string().contains("no authoritative baseline")); + assert!(err.to_string().contains("older than the required recovery fence")); } #[tokio::test] @@ -2646,6 +3203,11 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore { } objects.remove(&key); revisions.remove(&key); + drop(revisions); + drop(objects); + if self.error_after_commit_deletes.lock().await.remove(&key) { + return Err(EcstoreError::other("injected delete error after commit")); + } Ok(ObjectInfo::default()) } diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index e40a8fade..e54dac90c 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -12,6 +12,7 @@ for later deletion. ## Open Items +- `backlog-2102` rc.2/rc.3 empty scanner usage floor recovery: old DeleteBucket cleanup could synthesize an empty incomplete v2 usage primary/backup before leadership added an epoch, while newer scanners require a durable authoritative baseline identity. New scanners recognize only that exact serialized empty-fence shape, preserve its epoch through a CAS-protected recovery marker, and rebuild namespace coverage without treating zero usage as authoritative. Remove this recovery path and marker after rc.2 and rc.3 are no longer supported direct-upgrade sources. - `rustfs-6339` legacy bucket policy ID casing: earlier RustFS releases persisted the top-level policy identifier as "ID", while current writes use the S3-compatible "Id" spelling. Readers accept both spellings so retained bucket metadata remains usable after upgrade. Remove the legacy alias after migration tooling has rewritten every retained bucket policy using "ID". - `table-publication-fence-v1` table publication fencing: nodes that predate table and table-bucket publication fences can mutate live files while a new node is publishing a catalog pointer. New nodes retain exact object guards until the operator confirms that every serving node uses the new fences. Fleet confirmation also requires non-overlapping active warehouse prefixes and lifecycle workers that exclude table buckets. Remove the exact live-file fallback and the fleet-confirmation gate after the minimum supported RustFS release acquires table fences for registered-table mutations and table-bucket fences for unresolved-prefix mutations. - `table-catalog-strong-snapshot-v1` durable strong catalog snapshot compatibility: version 1 writes continue during mixed-version rollout until operators confirm that every serving node reads version 2, and version 1 table/view identifier collisions remain available only for cleanup. Remove version 1 writes and collision cleanup after the minimum supported RustFS release reads version 2 and every retained durable strong snapshot is collision-free and has been upgraded to version 2.