diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index d687c1081..ad2ddc3c2 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -552,18 +552,21 @@ pub(super) fn data_usage_info_has_persisted_baseline_identity(info: &DataUsageIn } pub(super) fn data_usage_info_is_bootstrap_pending(info: &DataUsageInfo) -> bool { - if info.last_update.is_none() || info.scanner_cycle.is_some() { + let Some(last_update) = info.last_update else { return false; - } + }; - let expected = DataUsageInfo { - last_update: info.last_update, - scanner_epoch: info.scanner_epoch, + info == &scanner_usage_bootstrap_marker(last_update, info.scanner_epoch) +} + +pub(super) fn scanner_usage_bootstrap_marker(last_update: std::time::SystemTime, scanner_epoch: Option) -> DataUsageInfo { + DataUsageInfo { + last_update: Some(last_update), + 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 { @@ -915,8 +918,8 @@ fn prepare_cycle_for_usage_floor_bootstrap( }, ) } - PersistedUsageFloorStartup::RecoveredLegacyEmptyFence => { - // The legacy empty fence proves only its leader epoch, not + PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence => { + // The legacy incomplete fence proves only its leader epoch, not // namespace coverage. Clear coverage while retaining the durable // cycle number so surviving caches cannot force a regression. let next = cycle_info.next; @@ -1438,6 +1441,7 @@ async fn fence_scanner_epoch_after_cycle_timeout( cycle_info: &mut CurrentCycle, cycle_revision: &mut DataUsageCacheRevision, leader_epoch: &mut u64, + allow_bootstrap_pending: bool, lock_lost: LockLost, ) -> bool where @@ -1451,7 +1455,7 @@ where cycle_info, cycle_revision, leader_epoch, - false, + allow_bootstrap_pending, ScannerCycleResetPolicy::None, ); tokio::pin!(claim); @@ -1473,6 +1477,7 @@ struct ScannerCycleDeadlineState<'a> { cycle_revision: &'a mut DataUsageCacheRevision, leader_epoch: &'a mut u64, cycle_budget: &'a ScannerCycleBudget, + allow_bootstrap_pending: bool, } fn cycle_timeout_requires_recovery(worker_stopped: bool, cycle_state_persisted: bool, generation_fenced: bool) -> bool { @@ -1494,6 +1499,7 @@ async fn handle_scanner_cycle_deadline( state.cycle_info, state.cycle_revision, state.leader_epoch, + state.allow_bootstrap_pending, guard.lock_lost_notified(), ) .await; @@ -2575,7 +2581,7 @@ async fn run_data_scanner_with_maintenance_state( match usage_floor_startup { PersistedUsageFloorStartup::Authoritative | PersistedUsageFloorStartup::BootstrapPending - | PersistedUsageFloorStartup::RecoveredLegacyEmptyFence => {} + | PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence => {} PersistedUsageFloorStartup::Missing => { if ctx.is_cancelled() || guard.is_lock_lost() { global_metrics().set_cycle(None).await; @@ -2670,8 +2676,8 @@ async fn run_data_scanner_with_maintenance_state( 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 + if usage_floor_startup == PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence + && let Err(err) = complete_legacy_incomplete_usage_floor_recovery(storeapi.clone(), leader_epoch).await { let error = err.to_string(); warn!( @@ -2744,6 +2750,7 @@ async fn run_data_scanner_with_maintenance_state( cycle_revision: &mut cycle_revision, leader_epoch: &mut leader_epoch, cycle_budget: &cycle_budget, + allow_bootstrap_pending: allow_usage_floor_bootstrap_pending, }, worker_stopped, &mut guard, @@ -3033,6 +3040,7 @@ async fn run_data_scanner_with_maintenance_state( cycle_revision: &mut cycle_revision, leader_epoch: &mut leader_epoch, cycle_budget: &cycle_budget, + allow_bootstrap_pending: allow_usage_floor_bootstrap_pending, }, worker_stopped, &mut guard, diff --git a/crates/scanner/src/scanner/cycle_state.rs b/crates/scanner/src/scanner/cycle_state.rs index 5788034ce..c35dbe036 100644 --- a/crates/scanner/src/scanner/cycle_state.rs +++ b/crates/scanner/src/scanner/cycle_state.rs @@ -26,7 +26,10 @@ 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"; +// Keep the published status value stable for operators that already alert on +// the empty-fence recovery introduced by backlog-2102. The same durable marker +// now also covers strictly validated data-bearing legacy fences. +const LEGACY_INCOMPLETE_USAGE_FLOOR_RECOVERY: &str = "legacy_empty_usage_floor"; const CACHE_CYCLE_AHEAD: &str = "cache_cycle_ahead"; const SCANNER_USAGE_STATE_RESET_MODE_FULL_REBUILD: &str = "full-rebuild"; @@ -182,9 +185,9 @@ pub(super) fn clear_scanner_cache_cycle_ahead() { } } -pub(super) fn record_legacy_empty_usage_floor_recovery_pending(leader_epoch: u64) { +pub(super) fn record_legacy_incomplete_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) + let same_recovery = previous.classification.as_deref() == Some(LEGACY_INCOMPLETE_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 { @@ -196,20 +199,20 @@ pub(super) fn record_legacy_empty_usage_floor_recovery_pending(leader_epoch: u64 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()), + classification: Some(LEGACY_INCOMPLETE_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()), + reason: Some("legacy incomplete 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) { +pub(super) fn clear_legacy_incomplete_usage_floor_recovery_status() { + if scanner_cycle_recovery_status().classification.as_deref() == Some(LEGACY_INCOMPLETE_USAGE_FLOOR_RECOVERY) { set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); } } @@ -1431,6 +1434,101 @@ async fn delete_usage_state_reset_slot( } } +#[derive(Clone, Copy)] +pub(super) enum ScannerUsageBootstrapPublishContext { + Initial, + Recovery, + Reset, +} + +enum ScannerUsageBootstrapPublishError { + Encode(serde_json::Error), + Reconcile(EcstoreError), + MissingEtag, + Save(EcstoreError), +} + +impl ScannerUsageBootstrapPublishError { + fn into_scanner_error(self, context: ScannerUsageBootstrapPublishContext) -> ScannerError { + let message = match context { + ScannerUsageBootstrapPublishContext::Initial => match self { + Self::Encode(err) => format!("failed to encode scanner usage baseline bootstrap: {err}"), + Self::Reconcile(err) => format!("failed to reconcile scanner usage bootstrap: {err}"), + Self::MissingEtag => "scanner usage bootstrap returned no ETag and could not be confirmed".to_string(), + Self::Save(err) => format!("failed to persist scanner usage bootstrap: {err}"), + }, + ScannerUsageBootstrapPublishContext::Recovery => match self { + Self::Encode(err) => format!("failed to encode recovered scanner usage bootstrap: {err}"), + Self::Reconcile(err) => format!("failed to reconcile recovered scanner usage bootstrap: {err}"), + Self::MissingEtag => "recovered scanner usage bootstrap returned no ETag and could not be confirmed".to_string(), + Self::Save(err) => format!("failed to recover legacy incomplete scanner usage floor: {err}"), + }, + ScannerUsageBootstrapPublishContext::Reset => match self { + Self::Encode(err) => format!("failed to encode scanner usage reset bootstrap marker: {err}"), + Self::Reconcile(err) => format!("failed to reconcile scanner usage reset bootstrap marker: {err}"), + Self::MissingEtag => "scanner usage reset bootstrap returned no ETag and could not be confirmed".to_string(), + Self::Save(err) if scanner_publication_epoch_changed(&err) => { + "scanner usage reset deferred by a movement epoch change".to_string() + } + Self::Save(EcstoreError::PreconditionFailed) => { + "scanner usage reset primary slot changed before bootstrap publish".to_string() + } + Self::Save(err) => format!("failed to persist scanner usage reset bootstrap: {err}"), + }, + }; + ScannerError::Other(message) + } +} + +pub(super) async fn publish_scanner_usage_bootstrap_primary( + storeapi: Arc, + expected_revision: &DataUsageCacheRevision, + expected_publication_epoch: u64, + leader_epoch: Option, + context: ScannerUsageBootstrapPublishContext, +) -> Result<(), ScannerError> { + async fn inner( + storeapi: Arc, + expected_revision: &DataUsageCacheRevision, + expected_publication_epoch: u64, + leader_epoch: Option, + ) -> Result<(), ScannerUsageBootstrapPublishError> { + let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::now(), leader_epoch); + let data = serde_json::to_vec(&marker).map_err(ScannerUsageBootstrapPublishError::Encode)?; + let save_result = save_config_with_publication_admission_for_epoch( + storeapi.clone(), + DATA_USAGE_OBJ_NAME_PATH.as_str(), + data.clone(), + expected_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()) + { + return Ok(()); + } + + let (persisted, revision) = read_config_with_revision(storeapi, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .map_err(ScannerUsageBootstrapPublishError::Reconcile)?; + if persisted.as_deref() == Some(data.as_slice()) && matches!(revision, DataUsageCacheRevision::Etag(_)) { + return Ok(()); + } + Err(match save_result { + Ok(_) => ScannerUsageBootstrapPublishError::MissingEtag, + Err(err) => ScannerUsageBootstrapPublishError::Save(err), + }) + } + + inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch) + .await + .map_err(|err| err.into_scanner_error(context)) +} + pub(super) async fn reset_scanner_usage_state_slots_for_full_rebuild( storeapi: Arc, slots: &[ScannerUsageStateResetSlot], @@ -1442,48 +1540,15 @@ pub(super) async fn reset_scanner_usage_state_slots_for_full_rebuild( .iter() .find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str()) .ok_or_else(|| ScannerError::Other("scanner usage reset primary slot was not inspected".to_string()))?; - let marker = DataUsageInfo { - last_update: Some(std::time::SystemTime::now()), - scanner_epoch: Some(leader_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 scanner usage reset bootstrap marker: {err}")))?; - let save_result = save_config_with_publication_admission_for_epoch( + publish_scanner_usage_bootstrap_primary( storeapi.clone(), - DATA_USAGE_OBJ_NAME_PATH.as_str(), - data.clone(), - primary.revision.preconditions(), + &primary.revision, expected_epoch, + Some(leader_epoch), + ScannerUsageBootstrapPublishContext::Reset, ) - .await; - if save_result - .as_ref() - .ok() - .and_then(|info| info.etag.as_deref()) - .is_some_and(|etag| !etag.is_empty()) - { - reset_paths.push(DATA_USAGE_OBJ_NAME_PATH.as_str().to_string()); - } else { - let (persisted, revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) - .await - .map_err(|err| ScannerError::Other(format!("failed to reconcile scanner usage reset bootstrap marker: {err}")))?; - if persisted.as_deref() != Some(data.as_slice()) || !matches!(revision, DataUsageCacheRevision::Etag(_)) { - return Err(ScannerError::Other(match save_result { - Ok(_) => "scanner usage reset bootstrap returned no ETag and could not be confirmed".to_string(), - Err(err) if scanner_publication_epoch_changed(&err) => { - "scanner usage reset deferred by a movement epoch change".to_string() - } - Err(EcstoreError::PreconditionFailed) => { - "scanner usage reset primary slot changed before bootstrap publish".to_string() - } - Err(err) => format!("failed to persist scanner usage reset bootstrap: {err}"), - })); - } - reset_paths.push(DATA_USAGE_OBJ_NAME_PATH.as_str().to_string()); - } + .await?; + reset_paths.push(DATA_USAGE_OBJ_NAME_PATH.as_str().to_string()); for slot in slots.iter().filter(|slot| slot.path != DATA_USAGE_OBJ_NAME_PATH.as_str()) { if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch).await? { @@ -1567,7 +1632,7 @@ pub async fn reset_scanner_usage_state_for_full_rebuild( } clear_scanner_usage_floor_failure(); - clear_legacy_empty_usage_floor_recovery_status(); + clear_legacy_incomplete_usage_floor_recovery_status(); set_scanner_cycle_recovery_status(recovery_status("healthy", None, false)); super::notify_scanner_cycle_recovery_wake(); info!( @@ -1614,33 +1679,48 @@ pub(super) enum PersistedUsageFloorStartup { Authoritative, Missing, BootstrapPending, - RecoveredLegacyEmptyFence, + RecoveredLegacyIncompleteFence, } #[derive(Clone, Debug)] -struct LegacyEmptyUsageFloorPrimary { +struct LegacyIncompleteUsageFloorPrimary { revision: DataUsageCacheRevision, epoch: u64, } #[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)] #[serde(deny_unknown_fields)] -struct LegacyEmptyUsageFloorRecoveryMarker { +struct LegacyIncompleteUsageFloorRecoveryMarker { schema_version: u16, primary_revision: String, leader_epoch: u64, } -async fn read_legacy_empty_usage_floor_recovery_marker( +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +struct LegacyIncompleteUsageFence { + claimable_epoch: Option, +} + +impl LegacyIncompleteUsageFence { + fn new(claimable_epoch: Option) -> Self { + Self { claimable_epoch } + } + + fn claimable_epoch(self) -> Option { + self.claimable_epoch + } +} + +async fn read_legacy_incomplete_usage_floor_recovery_marker( storeapi: Arc, -) -> Result, ScannerError> { +) -> 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) + 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())); @@ -1651,7 +1731,7 @@ async fn read_legacy_empty_usage_floor_recovery_marker( Ok(Some((marker, revision))) } -async fn clear_legacy_empty_usage_floor_recovery_marker( +async fn clear_legacy_incomplete_usage_floor_recovery_marker( storeapi: Arc, marker_revision: &DataUsageCacheRevision, expected_publication_epoch: u64, @@ -1685,11 +1765,11 @@ async fn clear_legacy_empty_usage_floor_recovery_marker( } } -pub(super) async fn complete_legacy_empty_usage_floor_recovery( +pub(super) async fn complete_legacy_incomplete_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 { + let Some((marker, marker_revision)) = read_legacy_incomplete_usage_floor_recovery_marker(storeapi.clone()).await? else { return Ok(()); }; if claimed_epoch <= marker.leader_epoch { @@ -1709,12 +1789,220 @@ pub(super) async fn complete_legacy_empty_usage_floor_recovery( 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(); + clear_legacy_incomplete_usage_floor_recovery_marker(storeapi, &marker_revision, expected_publication_epoch).await?; + clear_legacy_incomplete_usage_floor_recovery_status(); Ok(()) } -fn legacy_empty_usage_fence_epoch(data: &[u8], usage: &DataUsageInfo) -> Option> { +struct LegacyOptional { + present: bool, + value: Option, +} + +impl Default for LegacyOptional { + fn default() -> Self { + Self { + present: false, + value: None, + } + } +} + +fn deserialize_legacy_optional<'de, D, T>(deserializer: D) -> Result, D::Error> +where + D: serde::Deserializer<'de>, + T: Deserialize<'de>, +{ + Option::::deserialize(deserializer).map(|value| LegacyOptional { present: true, value }) +} + +struct LegacyUniqueMap(std::collections::HashMap); + +impl<'de, V> Deserialize<'de> for LegacyUniqueMap +where + V: Deserialize<'de>, +{ + fn deserialize(deserializer: D) -> Result + where + D: serde::Deserializer<'de>, + { + struct UniqueMapVisitor(std::marker::PhantomData); + + impl<'de, V> serde::de::Visitor<'de> for UniqueMapVisitor + where + V: Deserialize<'de>, + { + type Value = LegacyUniqueMap; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("a JSON object without duplicate keys") + } + + fn visit_map(self, mut entries: A) -> Result + where + A: serde::de::MapAccess<'de>, + { + let mut values = std::collections::HashMap::new(); + while let Some((key, value)) = entries.next_entry::()? { + match values.entry(key) { + std::collections::hash_map::Entry::Vacant(entry) => { + entry.insert(value); + } + std::collections::hash_map::Entry::Occupied(entry) => { + return Err(serde::de::Error::custom(format!("duplicate map key `{}`", entry.key()))); + } + } + } + Ok(LegacyUniqueMap(values)) + } + } + + deserializer.deserialize_map(UniqueMapVisitor(std::marker::PhantomData)) + } +} + +impl LegacyUniqueMap { + fn len(&self) -> usize { + self.0.len() + } + + fn get(&self, key: &str) -> Option<&V> { + self.0.get(key) + } + + fn iter(&self) -> impl Iterator { + self.0.iter() + } + + fn values(&self) -> impl Iterator { + self.0.values() + } +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyBucketTargetUsageWire { + #[serde(rename = "replication_pending_size")] + _replication_pending_size: serde::de::IgnoredAny, + #[serde(rename = "replication_failed_size")] + _replication_failed_size: serde::de::IgnoredAny, + #[serde(rename = "replicated_size")] + _replicated_size: serde::de::IgnoredAny, + #[serde(rename = "replica_size")] + _replica_size: serde::de::IgnoredAny, + #[serde(rename = "replication_pending_count")] + _replication_pending_count: serde::de::IgnoredAny, + #[serde(rename = "replication_failed_count")] + _replication_failed_count: serde::de::IgnoredAny, + #[serde(rename = "replicated_count")] + _replicated_count: serde::de::IgnoredAny, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyBucketUsageWire { + size: u64, + #[serde(rename = "replication_pending_size_v1")] + _replication_pending_size_v1: serde::de::IgnoredAny, + #[serde(rename = "replication_failed_size_v1")] + _replication_failed_size_v1: serde::de::IgnoredAny, + #[serde(rename = "replicated_size_v1")] + _replicated_size_v1: serde::de::IgnoredAny, + #[serde(rename = "replication_pending_count_v1")] + _replication_pending_count_v1: serde::de::IgnoredAny, + #[serde(rename = "replication_failed_count_v1")] + _replication_failed_count_v1: serde::de::IgnoredAny, + objects_count: u64, + #[serde(rename = "object_size_histogram")] + _object_size_histogram: LegacyUniqueMap, + #[serde(rename = "object_versions_histogram")] + _object_versions_histogram: LegacyUniqueMap, + versions_count: u64, + delete_markers_count: u64, + #[serde(rename = "replica_size")] + _replica_size: serde::de::IgnoredAny, + #[serde(rename = "replica_count")] + _replica_count: serde::de::IgnoredAny, + #[serde(rename = "replication_info")] + _replication_info: LegacyUniqueMap, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyDiskUsageStatusWire { + #[serde(rename = "disk_id")] + _disk_id: serde::de::IgnoredAny, + #[serde(rename = "pool_index")] + _pool_index: serde::de::IgnoredAny, + #[serde(rename = "set_index")] + _set_index: serde::de::IgnoredAny, + #[serde(rename = "disk_index")] + _disk_index: serde::de::IgnoredAny, + #[serde(rename = "last_update")] + _last_update: serde::de::IgnoredAny, + #[serde(rename = "snapshot_exists")] + _snapshot_exists: serde::de::IgnoredAny, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyTierStatsWire { + #[serde(rename = "total_size")] + _total_size: serde::de::IgnoredAny, + #[serde(rename = "num_versions")] + _num_versions: serde::de::IgnoredAny, + #[serde(rename = "num_objects")] + _num_objects: serde::de::IgnoredAny, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyAllTierStatsWire { + #[serde(rename = "tiers")] + _tiers: LegacyUniqueMap, +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct LegacyUsageWire { + #[serde(rename = "total_capacity")] + _total_capacity: serde::de::IgnoredAny, + #[serde(rename = "total_used_capacity")] + _total_used_capacity: serde::de::IgnoredAny, + #[serde(rename = "total_free_capacity")] + _total_free_capacity: serde::de::IgnoredAny, + #[serde(rename = "last_update")] + _last_update: serde::de::IgnoredAny, + #[serde(default, deserialize_with = "deserialize_legacy_optional")] + scanner_epoch: LegacyOptional, + objects_total_count: u64, + versions_total_count: u64, + delete_markers_total_count: u64, + objects_total_size: u64, + #[serde(rename = "replication_info")] + _replication_info: LegacyUniqueMap, + #[serde(default, deserialize_with = "deserialize_legacy_optional")] + tier_stats: LegacyOptional, + buckets_count: u64, + buckets_usage: LegacyUniqueMap, + usage_snapshot_complete: bool, + bucket_sizes: LegacyUniqueMap, + #[serde(rename = "disk_usage_status")] + _disk_usage_status: Vec, +} + +fn decode_legacy_usage_wire(data: &[u8], usage: &DataUsageInfo) -> Option { + let wire = serde_json::from_slice::(data).ok()?; + if wire.scanner_epoch.present != usage.scanner_epoch.is_some() + || wire.scanner_epoch.value != usage.scanner_epoch + || wire.tier_stats.present != usage.tier_stats.is_some() + { + return None; + } + Some(wire) +} + +fn legacy_empty_usage_fence(data: &[u8], usage: &DataUsageInfo) -> Option { if usage.last_update.is_none() || usage.scanner_cycle.is_some() { return None; } @@ -1730,55 +2018,88 @@ fn legacy_empty_usage_fence_epoch(data: &[u8], usage: &DataUsageInfo) -> Option< 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")) + let wire = decode_legacy_usage_wire(data, usage)?; + Some(LegacyIncompleteUsageFence::new(wire.scanner_epoch.value)) +} + +fn legacy_incomplete_usage_fence(data: &[u8], usage: &DataUsageInfo) -> Option { + legacy_empty_usage_fence(data, usage).or_else(|| legacy_non_empty_usage_fence(data, usage)) +} + +// RUSTFS_COMPAT_TODO(backlog-2181): accept rc.1-rc.3 usage floors that were fenced before a scanner cycle completed. Remove after those releases are no longer supported direct-upgrade sources. +fn legacy_non_empty_usage_fence(data: &[u8], usage: &DataUsageInfo) -> Option { + if usage.last_update.is_none() + || usage.scanner_cycle.is_some() + || usage.usage_snapshot_bootstrap_pending + || usage.usage_snapshot_complete + || usage.usage_snapshot_converged.is_some() + || usage.usage_snapshot_authoritative_baseline.is_some() + || !usage.usage_snapshot_set_states.is_empty() + || usage.usage_snapshot_partial + || usage.buckets_count == 0 + || u64::try_from(usage.buckets_usage.len()).ok() != Some(usage.buckets_count) { return None; } - Some(usage.scanner_epoch) + if usage.scanner_epoch.is_some_and(|epoch| epoch == 0 || epoch >= u64::MAX - 1) { + return None; + } + let wire = decode_legacy_usage_wire(data, usage)?; + if wire.usage_snapshot_complete + || wire.buckets_count == 0 + || u64::try_from(wire.buckets_usage.len()).ok() != Some(wire.buckets_count) + || wire.bucket_sizes.len() != wire.buckets_usage.len() + || wire + .buckets_usage + .iter() + .any(|(bucket, bucket_usage)| wire.bucket_sizes.get(bucket) != Some(&bucket_usage.size)) + { + return None; + } + let (objects, versions, delete_markers, size) = wire.buckets_usage.values().try_fold( + (0_u64, 0_u64, 0_u64, 0_u64), + |(objects, versions, delete_markers, size), bucket| { + Some(( + objects.checked_add(bucket.objects_count)?, + versions.checked_add(bucket.versions_count)?, + delete_markers.checked_add(bucket.delete_markers_count)?, + size.checked_add(bucket.size)?, + )) + }, + )?; + if (objects, versions, delete_markers, size) + != ( + wire.objects_total_count, + wire.versions_total_count, + wire.delete_markers_total_count, + wire.objects_total_size, + ) + { + return None; + } + Some(LegacyIncompleteUsageFence::new(wire.scanner_epoch.value)) } -async fn recover_legacy_empty_usage_floor( +async fn recover_legacy_incomplete_usage_floor( storeapi: Arc, - primary: LegacyEmptyUsageFloorPrimary, + primary: LegacyIncompleteUsageFloorPrimary, 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())); + return Err(ScannerError::Other("legacy incomplete scanner usage floor has no revision".to_string())); }; - let marker = LegacyEmptyUsageFloorRecoveryMarker { + let marker = LegacyIncompleteUsageFloorRecoveryMarker { 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? { + match read_legacy_incomplete_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(), + "scanner usage recovery marker conflicts with the persisted incomplete floor".to_string(), )); } Some(_) => {} @@ -1797,7 +2118,7 @@ async fn recover_legacy_empty_usage_floor( .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?; + let persisted = read_legacy_incomplete_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(), @@ -1808,52 +2129,26 @@ async fn recover_legacy_empty_usage_floor( } } - 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( + publish_scanner_usage_bootstrap_primary( storeapi.clone(), - DATA_USAGE_OBJ_NAME_PATH.as_str(), - data.clone(), - primary.revision.preconditions(), + &primary.revision, expected_publication_epoch, + Some(primary.epoch), + ScannerUsageBootstrapPublishContext::Recovery, ) - .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}"), - })) + .await?; + warn!( + target: "rustfs::scanner", + event = EVENT_SCANNER_PERSIST_STATE, + component = LOG_COMPONENT_SCANNER, + subsystem = LOG_SUBSYSTEM_RUNTIME, + // Keep the published state stable for existing empty-floor alerts. + state = "legacy_empty_usage_floor_recovered", + path = %DATA_USAGE_OBJ_NAME_PATH.as_str(), + scanner_epoch = primary.epoch, + "Scanner recovered a legacy incomplete usage floor" + ); + Ok(()) } pub(super) fn encode_scanner_cycle_state( @@ -2045,8 +2340,8 @@ fn resolve_bootstrap_backup_slot( if backup_epoch >= primary_epoch { update_persisted_usage_floor(resolution.floor, slot.usage, slot.backup_path)?; } - } else if let Some(epoch) = legacy_empty_usage_fence_epoch(slot.data, slot.usage) { - if let Some(epoch) = epoch { + } else if let Some(fence) = legacy_incomplete_usage_fence(slot.data, slot.usage) { + if let Some(epoch) = fence.claimable_epoch() { resolution.floor.leader_epoch = resolution.floor.leader_epoch.max(epoch); } } else { @@ -2056,10 +2351,12 @@ fn resolve_bootstrap_backup_slot( } return Ok(BootstrapBackupAction::Resume); } - let compatible_empty_fence = legacy_empty_usage_fence_epoch(slot.data, slot.usage).is_some_and(|epoch| { - epoch.is_none_or(|epoch| slot.bootstrap_epoch.is_some_and(|bootstrap_epoch| epoch <= bootstrap_epoch)) + let compatible_incomplete_fence = legacy_incomplete_usage_fence(slot.data, slot.usage).is_some_and(|fence| { + fence + .claimable_epoch() + .is_none_or(|epoch| slot.bootstrap_epoch.is_some_and(|bootstrap_epoch| epoch <= bootstrap_epoch)) }); - if compatible_empty_fence { + if compatible_incomplete_fence { return Ok(BootstrapBackupAction::Resume); } if slot.recovered_bootstrap && data_usage_info_has_persisted_baseline_identity(slot.usage) { @@ -2084,7 +2381,7 @@ 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 recovery_marker = read_legacy_incomplete_usage_floor_recovery_marker(storeapi.clone()).await?; let mut floor = PersistedUsageFloor::default(); let mut found_any = false; let mut bootstrap_pending = false; @@ -2099,7 +2396,7 @@ pub(super) async fn persisted_usage_floor_for_startup( 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 mut legacy_incomplete_primary: Option = None; 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(); @@ -2148,14 +2445,12 @@ pub(super) async fn persisted_usage_floor_for_startup( } 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()); + if let Some(fence) = legacy_incomplete_usage_fence(&data, &usage) { + if is_v2_path && let Some(epoch) = fence.claimable_epoch() { + legacy_incomplete_primary = Some(LegacyIncompleteUsageFloorPrimary { revision, epoch }); } + } else { + unrecoverable_baseline_path.get_or_insert_with(|| primary_path.to_string()); } None } else { @@ -2243,7 +2538,7 @@ pub(super) async fn persisted_usage_floor_for_startup( 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() { + if legacy_incomplete_usage_fence(&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 @@ -2294,17 +2589,18 @@ pub(super) async fn persisted_usage_floor_for_startup( if allow_missing_for_bootstrap && unrecoverable_baseline_path.is_none() && stale_authoritative_path.is_none() - && let Some(mut primary) = legacy_empty_primary + && let Some(mut primary) = legacy_incomplete_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); + let leader_epoch = primary.epoch; + recover_legacy_incomplete_usage_floor(storeapi.clone(), primary, read_epoch).await?; + record_legacy_incomplete_usage_floor_recovery_pending(leader_epoch); return Ok(( PersistedUsageFloor { next_cycle: 0, - leader_epoch: primary.epoch, + leader_epoch, }, - PersistedUsageFloorStartup::RecoveredLegacyEmptyFence, + PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence, )); } if let Some(path) = stale_authoritative_path { @@ -2317,7 +2613,7 @@ pub(super) async fn persisted_usage_floor_for_startup( } 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" + "persisted scanner usage floor from {path} has no authoritative baseline or newer valid backup; recover with POST /rustfs/admin/v3/scanner/usage-state/reset using mode full-rebuild" ))); } if !allow_missing_for_bootstrap { @@ -2369,8 +2665,8 @@ pub(super) async fn persisted_usage_floor_for_startup( .as_ref() .map(|(marker, _)| marker.leader_epoch) .unwrap_or(floor.leader_epoch); - record_legacy_empty_usage_floor_recovery_pending(recovery_epoch); - PersistedUsageFloorStartup::RecoveredLegacyEmptyFence + record_legacy_incomplete_usage_floor_recovery_pending(recovery_epoch); + PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence } else if bootstrap_pending { if let Some(path) = unrecoverable_baseline_path { return Err(ScannerError::Other(format!( @@ -2384,7 +2680,7 @@ pub(super) async fn persisted_usage_floor_for_startup( 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 { + match clear_legacy_incomplete_usage_floor_recovery_marker(storeapi.clone(), marker_revision, read_epoch).await { Ok(()) => true, Err(err) => { warn!( @@ -2407,7 +2703,7 @@ pub(super) async fn persisted_usage_floor_for_startup( )); }; if marker_cleared { - clear_legacy_empty_usage_floor_recovery_status(); + clear_legacy_incomplete_usage_floor_recovery_status(); } clear_scanner_usage_floor_failure(); return Ok((floor, state)); diff --git a/crates/scanner/src/scanner/leadership.rs b/crates/scanner/src/scanner/leadership.rs index d273d133a..8283a987a 100644 --- a/crates/scanner/src/scanner/leadership.rs +++ b/crates/scanner/src/scanner/leadership.rs @@ -185,42 +185,14 @@ pub(super) async fn initialize_usage_baseline_bootstrap( "scanner usage baseline bootstrap is blocked by data movement".to_string(), )); }; - let baseline = DataUsageInfo { - last_update: Some(std::time::SystemTime::now()), - usage_snapshot_converged: Some(false), - usage_snapshot_bootstrap_pending: true, - ..Default::default() - }; - let data = serde_json::to_vec(&baseline) - .map_err(|err| ScannerError::Other(format!("failed to encode scanner usage baseline bootstrap: {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(), + publish_scanner_usage_bootstrap_primary( + storeapi, + &DataUsageCacheRevision::Missing, expected_epoch, + None, + ScannerUsageBootstrapPublishContext::Initial, ) - .await; - if save_result - .as_ref() - .ok() - .and_then(|info| info.etag.as_deref()) - .is_some_and(|etag| !etag.is_empty()) - { - 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 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(_) => "scanner usage bootstrap returned no ETag and could not be confirmed".to_string(), - Err(err) => format!("failed to persist scanner usage bootstrap: {err}"), - })) + .await } pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch( diff --git a/crates/scanner/src/scanner/tests.rs b/crates/scanner/src/scanner/tests.rs index 263501089..e4c836c2e 100644 --- a/crates/scanner/src/scanner/tests.rs +++ b/crates/scanner/src/scanner/tests.rs @@ -462,6 +462,7 @@ async fn cycle_budget_persist_cursor_failure_is_recovery_required() { &mut cycle, &mut revision, &mut leader_epoch, + false, std::future::pending(), ) .await; @@ -478,6 +479,52 @@ async fn cycle_budget_persist_cursor_failure_is_recovery_required() { assert!(report.leader_lease_without_progress); } +#[tokio::test] +async fn cycle_budget_fence_accepts_bootstrap_pending_usage_marker() { + let store = Arc::new(MemoryConfigStore::default()); + initialize_usage_baseline_bootstrap(store.clone()) + .await + .expect("usage reset should publish a bootstrap marker"); + + let ctx = CancellationToken::new(); + let mut revision = DataUsageCacheRevision::Missing; + let mut cycle = CurrentCycle { + current: 12, + next: 12, + ..Default::default() + }; + let mut leader_epoch = 0; + + let fenced = fence_scanner_epoch_after_cycle_timeout( + &ctx, + store.clone(), + &mut cycle, + &mut revision, + &mut leader_epoch, + true, + std::future::pending(), + ) + .await; + + assert!(fenced, "a valid reset bootstrap marker must not force cycle recovery after budget expiry"); + assert!(!cycle_timeout_requires_recovery(true, true, fenced)); + assert_eq!(leader_epoch, 1); + + let persisted_cycle = read_config(store.clone(), &DATA_USAGE_BLOOM_NAME_PATH) + .await + .expect("timeout fence should persist the next leader epoch"); + let (_, persisted_epoch) = decode_scanner_cycle_state(&persisted_cycle).expect("persisted epoch fence should decode"); + assert_eq!(persisted_epoch, 1); + + let usage = read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("timeout fence should keep the bootstrap usage marker"); + let usage = serde_json::from_slice::(&usage).expect("bootstrap marker should decode"); + assert!(data_usage_info_is_bootstrap_pending(&usage)); + assert_eq!(usage.scanner_epoch, Some(1)); + assert!(!data_usage_info_has_persisted_baseline_identity(&usage)); +} + #[tokio::test] async fn cycle_budget_deadline_handler_fences_and_releases_guard() { let (_temp_dir, store) = setup_scanner_cycle_store().await; @@ -515,6 +562,7 @@ async fn cycle_budget_deadline_handler_fences_and_releases_guard() { cycle_revision: &mut cycle_revision, leader_epoch: &mut leader_epoch, cycle_budget: &budget, + allow_bootstrap_pending: false, }, true, &mut guard, @@ -2237,6 +2285,53 @@ fn rc3_legacy_empty_usage_fence(epoch: Option) -> Vec { serde_json::to_vec(&value).expect("rc.3 legacy empty usage fence fixture should encode") } +fn rc3_legacy_non_empty_usage_fence(epoch: Option) -> Vec { + // Pinned rc.3 field set. Leadership preserved this data and added only + // scanner_epoch when the producing scanner cycle had not completed. + const RC3_NON_EMPTY_USAGE_FENCE: &str = r#"{ + "total_capacity":2000000000, + "total_used_capacity":1000000000, + "total_free_capacity":1000000000, + "last_update":{"secs_since_epoch":1,"nanos_since_epoch":0}, + "objects_total_count":156382067, + "versions_total_count":156382070, + "delete_markers_total_count":3, + "objects_total_size":987654321, + "replication_info":{}, + "buckets_count":1, + "buckets_usage":{ + "photos":{ + "size":987654321, + "replication_pending_size_v1":0, + "replication_failed_size_v1":0, + "replicated_size_v1":0, + "replication_pending_count_v1":0, + "replication_failed_count_v1":0, + "objects_count":156382067, + "object_size_histogram":{}, + "object_versions_histogram":{}, + "versions_count":156382070, + "delete_markers_count":3, + "replica_size":0, + "replica_count":0, + "replication_info":{} + } + }, + "usage_snapshot_complete":false, + "bucket_sizes":{"photos":987654321}, + "disk_usage_status":[] + }"#; + let mut value = serde_json::from_str::(RC3_NON_EMPTY_USAGE_FENCE) + .expect("pinned rc.3 non-empty usage fence should decode"); + if let Some(epoch) = epoch { + value + .as_object_mut() + .expect("legacy non-empty usage fence should be a JSON object") + .insert("scanner_epoch".to_string(), serde_json::Value::from(epoch)); + } + serde_json::to_vec(&value).expect("rc.3 legacy non-empty usage fence fixture should encode") +} + #[tokio::test] async fn scanner_usage_floor_recovers_rc3_empty_fences_and_preserves_cycle_number() { let store = Arc::new(MemoryConfigStore::default()); @@ -2265,7 +2360,7 @@ async fn scanner_usage_floor_recovers_rc3_empty_fences_and_preserves_cycle_numbe leader_epoch: 7, } ); - assert_eq!(startup, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + assert_eq!(startup, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); let primary = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) .await @@ -2279,7 +2374,7 @@ async fn scanner_usage_floor_recovers_rc3_empty_fences_and_preserves_cycle_numbe .await .expect("recovery marker should survive a restart before leadership claim"); assert_eq!(restart_floor.leader_epoch, 7); - assert_eq!(restart_state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + assert_eq!(restart_state, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); let mut cycle = CurrentCycle { current: 17_117, next: 17_118, @@ -2368,7 +2463,7 @@ async fn scanner_usage_floor_recovers_rc3_empty_fences_and_preserves_cycle_numbe assert!(!persisted_reset.info.snapshot_complete); assert!(persisted_reset.cache.is_empty()); } - complete_legacy_empty_usage_floor_recovery(store.clone(), leader_epoch) + complete_legacy_incomplete_usage_floor_recovery(store.clone(), leader_epoch) .await .expect("leadership claim should retire the recovery marker"); assert!(matches!( @@ -2382,6 +2477,277 @@ async fn scanner_usage_floor_recovers_rc3_empty_fences_and_preserves_cycle_numbe assert_eq!(claimed_state, PersistedUsageFloorStartup::BootstrapPending); } +#[tokio::test] +async fn scanner_usage_floor_recovers_rc3_non_empty_incomplete_fence() { + 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_non_empty_usage_fence(Some(13))); + store.revisions.lock().await.insert(primary_key, 1); + + save_config( + store.clone(), + LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(), + rc3_legacy_non_empty_usage_fence(None), + ) + .await + .expect("legacy usage should persist"); + + let (floor, startup) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("rc.3 non-empty incomplete fence should enter recovery"); + assert_eq!( + floor, + PersistedUsageFloor { + next_cycle: 0, + leader_epoch: 13, + } + ); + assert_eq!(startup, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); + + let primary = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("recovered usage bootstrap should replace the old floor"); + let pending = serde_json::from_slice::(&primary).expect("recovered usage bootstrap should decode"); + assert!(data_usage_info_is_bootstrap_pending(&pending)); + assert!(!data_usage_info_has_persisted_baseline_identity(&pending)); + assert_eq!(pending.scanner_epoch, Some(13)); + assert!(read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await.is_ok()); +} + +#[tokio::test] +async fn scanner_usage_floor_prefers_newer_backup_over_rc3_non_empty_incomplete_fence() { + 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, rc3_legacy_non_empty_usage_fence(Some(13))); + + let mut backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 2); + backup.scanner_epoch = Some(14); + backup.scanner_cycle = Some(9845); + save_config( + store.clone(), + &format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()), + serde_json::to_vec(&backup).expect("newer backup should encode"), + ) + .await + .expect("newer backup should persist"); + + let (floor, startup) = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect("newer authoritative backup should win over the old incomplete floor"); + assert_eq!( + floor, + PersistedUsageFloor { + next_cycle: 9846, + leader_epoch: 14, + } + ); + assert_eq!(startup, PersistedUsageFloorStartup::Authoritative); + assert!(matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + +#[tokio::test] +async fn scanner_usage_floor_rejects_noncanonical_non_empty_incomplete_fences() { + let base = serde_json::from_slice::(&rc3_legacy_non_empty_usage_fence(Some(13))) + .expect("pinned rc.3 usage fence should decode"); + let mut cases = Vec::new(); + + let mut unknown_top_level = base.clone(); + unknown_top_level["future_field"] = serde_json::Value::Bool(true); + cases.push(("unknown top-level field", unknown_top_level)); + + let mut unknown_bucket_field = base.clone(); + unknown_bucket_field["buckets_usage"]["photos"]["future_field"] = serde_json::Value::Bool(true); + cases.push(("unknown bucket field", unknown_bucket_field)); + + let mut wrong_bucket_size = base.clone(); + wrong_bucket_size["bucket_sizes"]["photos"] = serde_json::Value::from(987_654_320_u64); + cases.push(("bucket size mismatch", wrong_bucket_size)); + + let mut wrong_total = base.clone(); + wrong_total["objects_total_count"] = serde_json::Value::from(156_382_068_u64); + cases.push(("object total mismatch", wrong_total)); + + let mut wrong_versions = base.clone(); + wrong_versions["versions_total_count"] = serde_json::Value::from(156_382_071_u64); + cases.push(("version total mismatch", wrong_versions)); + + let mut wrong_delete_markers = base.clone(); + wrong_delete_markers["delete_markers_total_count"] = serde_json::Value::from(4_u64); + cases.push(("delete marker total mismatch", wrong_delete_markers)); + + let mut wrong_total_size = base.clone(); + wrong_total_size["objects_total_size"] = serde_json::Value::from(987_654_320_u64); + cases.push(("object size total mismatch", wrong_total_size)); + + let mut wrong_cardinality = base.clone(); + wrong_cardinality["buckets_count"] = serde_json::Value::from(2_u64); + cases.push(("bucket cardinality mismatch", wrong_cardinality)); + + let mut overflow = base.clone(); + let mut overflow_bucket = overflow["buckets_usage"]["photos"].clone(); + overflow_bucket["objects_count"] = serde_json::Value::from(u64::MAX); + overflow_bucket["size"] = serde_json::Value::from(0_u64); + overflow["buckets_usage"]["overflow"] = overflow_bucket; + overflow["bucket_sizes"]["overflow"] = serde_json::Value::from(0_u64); + overflow["buckets_count"] = serde_json::Value::from(2_u64); + cases.push(("checked total overflow", overflow)); + + let mut invalid_epoch = base; + invalid_epoch["scanner_epoch"] = serde_json::Value::from(0_u64); + cases.push(("invalid epoch", invalid_epoch)); + + for (case, value) in cases { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + let original = serde_json::to_vec(&value).expect("noncanonical usage fixture should encode"); + store.objects.lock().await.insert(primary_key.clone(), original.clone()); + store.revisions.lock().await.insert(primary_key, 1); + + let err = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("noncanonical incomplete usage must remain fail-closed"); + assert!(err.to_string().contains("usage-state/reset"), "unexpected error for {case}: {err}"); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("rejected usage primary should remain"), + original, + "rejected primary changed for {case}" + ); + assert!( + matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + ), + "recovery marker should not be written for {case}" + ); + } +} + +#[tokio::test] +async fn scanner_usage_floor_rejects_duplicate_legacy_fields() { + let base = String::from_utf8(rc3_legacy_non_empty_usage_fence(Some(13))).expect("pinned rc.3 usage fence should be UTF-8"); + let base_value = serde_json::from_str::(&base).expect("pinned rc.3 usage fence should decode"); + let duplicate_top_level = base.replacen( + "\"objects_total_count\":156382067", + "\"objects_total_count\":156382067,\"objects_total_count\":156382067", + 1, + ); + let duplicate_bucket = base.replacen("\"size\":987654321", "\"size\":987654321,\"size\":987654321", 1); + let bucket = serde_json::to_string(&base_value["buckets_usage"]["photos"]).expect("pinned rc.3 bucket usage should encode"); + let bucket_map = format!("\"buckets_usage\":{{\"photos\":{bucket}}}"); + let duplicate_bucket_key = + base.replacen(&bucket_map, &format!("\"buckets_usage\":{{\"photos\":{bucket},\"photos\":{bucket}}}"), 1); + let duplicate_bucket_size_key = base.replacen( + "\"bucket_sizes\":{\"photos\":987654321}", + "\"bucket_sizes\":{\"photos\":987654321,\"photos\":987654321}", + 1, + ); + let duplicate_histogram_key = + base.replacen("\"object_size_histogram\":{}", "\"object_size_histogram\":{\"small\":1,\"small\":1}", 1); + let mut duplicate_target_field = base.clone(); + let target_map = "\"replication_info\":{\"target\":{\"replication_pending_size\":0,\"replication_failed_size\":0,\"replicated_size\":0,\"replica_size\":0,\"replication_pending_count\":0,\"replication_failed_count\":0,\"replicated_count\":0,\"replicated_count\":0}}"; + let target_offset = duplicate_target_field + .rfind("\"replication_info\":{}") + .expect("pinned fixture should contain bucket replication info"); + duplicate_target_field.replace_range(target_offset..target_offset + "\"replication_info\":{}".len(), target_map); + let target = "{\"replication_pending_size\":0,\"replication_failed_size\":0,\"replicated_size\":0,\"replica_size\":0,\"replication_pending_count\":0,\"replication_failed_count\":0,\"replicated_count\":0}"; + let mut duplicate_target_key = base.clone(); + let target_offset = duplicate_target_key + .rfind("\"replication_info\":{}") + .expect("pinned fixture should contain bucket replication info"); + let target_map = format!("\"replication_info\":{{\"target\":{target},\"target\":{target}}}"); + duplicate_target_key.replace_range(target_offset..target_offset + "\"replication_info\":{}".len(), &target_map); + + let tier = "{\"total_size\":0,\"num_versions\":0,\"num_objects\":0}"; + let duplicate_tier_key = base.replacen( + '{', + &format!("{{\"tier_stats\":{{\"tiers\":{{\"STANDARD\":{tier},\"STANDARD\":{tier}}}}},"), + 1, + ); + + for (case, original, expected_error) in [ + ("top-level field", duplicate_top_level.into_bytes(), "duplicate field"), + ("bucket field", duplicate_bucket.into_bytes(), "duplicate field"), + ("replication target field", duplicate_target_field.into_bytes(), "duplicate field"), + ("bucket map key", duplicate_bucket_key.into_bytes(), "usage-state/reset"), + ("bucket size map key", duplicate_bucket_size_key.into_bytes(), "usage-state/reset"), + ("histogram map key", duplicate_histogram_key.into_bytes(), "usage-state/reset"), + ("replication target map key", duplicate_target_key.into_bytes(), "usage-state/reset"), + ("tier map key", duplicate_tier_key.into_bytes(), "usage-state/reset"), + ] { + 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(), original.clone()); + store.revisions.lock().await.insert(primary_key, 1); + + let err = match persisted_usage_floor_for_startup(store.clone(), true).await { + Err(err) => err, + Ok(result) => panic!("duplicate {case} must remain fail-closed: {result:?}"), + }; + assert!(err.to_string().contains(expected_error), "unexpected error for {case}: {err}"); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("rejected duplicate-field primary should remain"), + original, + "rejected primary changed for {case}" + ); + assert!( + matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + ), + "recovery marker should not be written for {case}" + ); + } +} + +#[tokio::test] +async fn scanner_usage_floor_rejects_current_schema_incomplete_snapshot() { + let store = Arc::new(MemoryConfigStore::default()); + let primary_key = memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_OBJ_NAME_PATH.as_str()); + let mut current = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 2); + current.usage_snapshot_complete = false; + current.scanner_epoch = Some(13); + current.scanner_cycle = None; + let original = serde_json::to_vec(¤t).expect("current incomplete usage should encode"); + assert!( + serde_json::from_slice::(&original) + .expect("current incomplete usage should decode") + .get("usage_snapshot_partial") + .is_some() + ); + store.objects.lock().await.insert(primary_key.clone(), original.clone()); + store.revisions.lock().await.insert(primary_key, 1); + + let err = persisted_usage_floor_for_startup(store.clone(), true) + .await + .expect_err("current schema incomplete usage must remain fail-closed"); + assert!(err.to_string().contains("usage-state/reset"), "unexpected error: {err}"); + assert_eq!( + read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()) + .await + .expect("rejected current usage primary should remain"), + original + ); + assert!(matches!( + read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await, + Err(EcstoreError::ConfigNotFound) + )); +} + #[tokio::test] async fn scanner_usage_floor_recovery_preserves_newer_authoritative_companion_floor() { for companion_path in [ @@ -2430,7 +2796,7 @@ async fn scanner_usage_floor_recovery_preserves_newer_authoritative_companion_fl .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); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); } } @@ -2485,7 +2851,7 @@ async fn scanner_usage_floor_recovery_fences_non_authoritative_legacy_backup() { .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); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); } } } @@ -2514,7 +2880,7 @@ async fn scanner_usage_floor_recovery_resumes_after_marker_only_crash_point() { .await .expect("the durable marker should resume the primary conversion"); assert_eq!(floor.leader_epoch, 7); - assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); let recovered = read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str()) .await .expect("recovered bootstrap should replace the legacy primary"); @@ -2540,7 +2906,7 @@ async fn scanner_usage_floor_recovery_reconciles_marker_post_commit_error() { .await .expect("a committed recovery marker should reconcile after an ambiguous error"); assert_eq!(floor.leader_epoch, 7); - assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); assert!(read_config(store, DATA_USAGE_RECOVERY_PATH.as_str()).await.is_ok()); } @@ -2579,7 +2945,7 @@ async fn scanner_usage_floor_recovery_reconciles_marker_delete_post_commit_error .await .insert(memory_config_key(RUSTFS_META_BUCKET, DATA_USAGE_RECOVERY_PATH.as_str())); - complete_legacy_empty_usage_floor_recovery(store.clone(), leader_epoch) + complete_legacy_incomplete_usage_floor_recovery(store.clone(), leader_epoch) .await .expect("a committed marker delete should reconcile after an ambiguous error"); assert!(matches!( @@ -2621,12 +2987,12 @@ async fn scanner_usage_floor_recovery_retry_budget_uses_marker_epoch_identity() .await .expect("claimed bootstrap should retain its recovery identity"); assert_eq!(floor.leader_epoch, 8); - assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyEmptyFence); + assert_eq!(state, PersistedUsageFloorStartup::RecoveredLegacyIncompleteFence); 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(); + clear_legacy_incomplete_usage_floor_recovery_status(); } #[tokio::test] @@ -2782,7 +3148,7 @@ fn scanner_usage_floor_failure_is_exposed_and_cleared() { #[test] #[serial] fn scanner_usage_floor_recovery_stays_retryable_until_claim_cleanup() { - record_legacy_empty_usage_floor_recovery_pending(7); + record_legacy_incomplete_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")); @@ -2792,12 +3158,12 @@ fn scanner_usage_floor_recovery_stays_retryable_until_claim_cleanup() { let first_detected = pending.first_detected_at_unix_secs; assert!(record_scanner_cycle_recovery_retry(3)); - record_legacy_empty_usage_floor_recovery_pending(7); + record_legacy_incomplete_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(); + clear_legacy_incomplete_usage_floor_recovery_status(); assert_eq!(scanner_cycle_recovery_status().state, "healthy"); } diff --git a/docs/architecture/compat-cleanup-register.md b/docs/architecture/compat-cleanup-register.md index d598bd553..2fa6f1319 100644 --- a/docs/architecture/compat-cleanup-register.md +++ b/docs/architecture/compat-cleanup-register.md @@ -13,6 +13,7 @@ - `tokio-tar-extension-limits` bounded archive parser hardening: Snowball extraction depends on per-entry and cumulative GNU long-name, GNU long-link, and PAX extension limits; physical-entry, GNU sparse-map, and sparse-continuation limits; cancellation-safe sparse parsing; and fused entry streams after parser errors. The released tokio-tar API does not provide this complete boundary. Keep the reviewed fork pin until astral-sh/tokio-tar#118 is merged and one published tokio-tar release contains every listed capability with the Snowball regression fixtures passing against that release. - `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. +- `backlog-2181` rc.1-rc.3 non-empty scanner usage floor recovery: leadership fencing in those releases can stamp scanner_epoch onto a real bucket-usage snapshot before any scanner cycle completed, leaving a non-empty floor with no scanner_cycle and no authoritative baseline identity. New scanners recognize only this consistent incomplete fenced shape, preserve the epoch through the CAS-protected recovery marker, and rebuild namespace coverage without treating the old usage data as authoritative. Remove this recovery path after rc.1, rc.2, and rc.3 are no longer supported direct-upgrade sources. - `s3gate-metadata-xml` persisted bucket XML migration: mixed-version site-replication peers, retained `.metadata.bin` objects, and backup archives can all carry XML written by the s3s codec, so the gateway migration must keep the legacy codec available until every stored form has crossed a verified rewrite boundary. Remove the legacy s3s parser and serializer only after the minimum supported direct-upgrade release reads and writes every persisted XML configuration family through the gateway codec, every supported mixed-version site-replication topology has completed its writer upgrade, and migration tooling has verified or rewritten every retained bucket metadata object and restorable backup archive. - `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.