fix(scanner): recover fenced incomplete usage floors (#7055)

* fix(scanner): recover fenced incomplete usage floors

* fix(scanner): validate legacy usage floor shape

Co-Authored-By: heihutu <heihutu@gmail.com>

---------

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-09-02 20:32:31 +08:00
committed by GitHub
parent afc66b7182
commit 99f85ca2b1
5 changed files with 861 additions and 218 deletions
+21 -13
View File
@@ -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<u64>) -> 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<Store, LockLost>(
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<Store>(
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,
+454 -158
View File
@@ -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<impl ScannerObjectIO + ScannerConfigObjectDelete>,
expected_revision: &DataUsageCacheRevision,
expected_publication_epoch: u64,
leader_epoch: Option<u64>,
context: ScannerUsageBootstrapPublishContext,
) -> Result<(), ScannerError> {
async fn inner(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
expected_revision: &DataUsageCacheRevision,
expected_publication_epoch: u64,
leader_epoch: Option<u64>,
) -> 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<impl ScannerObjectIO + ScannerConfigObjectDelete>,
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<u64>,
}
impl LegacyIncompleteUsageFence {
fn new(claimable_epoch: Option<u64>) -> Self {
Self { claimable_epoch }
}
fn claimable_epoch(self) -> Option<u64> {
self.claimable_epoch
}
}
async fn read_legacy_incomplete_usage_floor_recovery_marker(
storeapi: Arc<impl ScannerObjectIO>,
) -> Result<Option<(LegacyEmptyUsageFloorRecoveryMarker, DataUsageCacheRevision)>, ScannerError> {
) -> Result<Option<(LegacyIncompleteUsageFloorRecoveryMarker, DataUsageCacheRevision)>, 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::<LegacyEmptyUsageFloorRecoveryMarker>(&data)
let marker = serde_json::from_slice::<LegacyIncompleteUsageFloorRecoveryMarker>(&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<impl ScannerObjectIO + ScannerConfigObjectDelete>,
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<impl ScannerObjectIO + ScannerConfigObjectDelete>,
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<Option<u64>> {
struct LegacyOptional<T> {
present: bool,
value: Option<T>,
}
impl<T> Default for LegacyOptional<T> {
fn default() -> Self {
Self {
present: false,
value: None,
}
}
}
fn deserialize_legacy_optional<'de, D, T>(deserializer: D) -> Result<LegacyOptional<T>, D::Error>
where
D: serde::Deserializer<'de>,
T: Deserialize<'de>,
{
Option::<T>::deserialize(deserializer).map(|value| LegacyOptional { present: true, value })
}
struct LegacyUniqueMap<V>(std::collections::HashMap<String, V>);
impl<'de, V> Deserialize<'de> for LegacyUniqueMap<V>
where
V: Deserialize<'de>,
{
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
struct UniqueMapVisitor<V>(std::marker::PhantomData<V>);
impl<'de, V> serde::de::Visitor<'de> for UniqueMapVisitor<V>
where
V: Deserialize<'de>,
{
type Value = LegacyUniqueMap<V>;
fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("a JSON object without duplicate keys")
}
fn visit_map<A>(self, mut entries: A) -> Result<Self::Value, A::Error>
where
A: serde::de::MapAccess<'de>,
{
let mut values = std::collections::HashMap::new();
while let Some((key, value)) = entries.next_entry::<String, V>()? {
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<V> LegacyUniqueMap<V> {
fn len(&self) -> usize {
self.0.len()
}
fn get(&self, key: &str) -> Option<&V> {
self.0.get(key)
}
fn iter(&self) -> impl Iterator<Item = (&String, &V)> {
self.0.iter()
}
fn values(&self) -> impl Iterator<Item = &V> {
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<u64>,
#[serde(rename = "object_versions_histogram")]
_object_versions_histogram: LegacyUniqueMap<u64>,
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<LegacyBucketTargetUsageWire>,
}
#[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<LegacyTierStatsWire>,
}
#[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<u64>,
objects_total_count: u64,
versions_total_count: u64,
delete_markers_total_count: u64,
objects_total_size: u64,
#[serde(rename = "replication_info")]
_replication_info: LegacyUniqueMap<LegacyBucketTargetUsageWire>,
#[serde(default, deserialize_with = "deserialize_legacy_optional")]
tier_stats: LegacyOptional<LegacyAllTierStatsWire>,
buckets_count: u64,
buckets_usage: LegacyUniqueMap<LegacyBucketUsageWire>,
usage_snapshot_complete: bool,
bucket_sizes: LegacyUniqueMap<u64>,
#[serde(rename = "disk_usage_status")]
_disk_usage_status: Vec<LegacyDiskUsageStatusWire>,
}
fn decode_legacy_usage_wire(data: &[u8], usage: &DataUsageInfo) -> Option<LegacyUsageWire> {
let wire = serde_json::from_slice::<LegacyUsageWire>(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<LegacyIncompleteUsageFence> {
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::<serde_json::Value>(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<LegacyIncompleteUsageFence> {
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<LegacyIncompleteUsageFence> {
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<impl ScannerObjectIO + ScannerConfigObjectDelete>,
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<String> = None;
let mut stale_authoritative_path: Option<String> = None;
let mut legacy_empty_primary: Option<LegacyEmptyUsageFloorPrimary> = None;
let mut legacy_incomplete_primary: Option<LegacyIncompleteUsageFloorPrimary> = 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));
+6 -34
View File
@@ -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(
+379 -13
View File
@@ -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::<DataUsageInfo>(&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<u64>) -> Vec<u8> {
serde_json::to_vec(&value).expect("rc.3 legacy empty usage fence fixture should encode")
}
fn rc3_legacy_non_empty_usage_fence(epoch: Option<u64>) -> Vec<u8> {
// 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::<serde_json::Value>(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::<DataUsageInfo>(&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::<serde_json::Value>(&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::<serde_json::Value>(&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(&current).expect("current incomplete usage should encode");
assert!(
serde_json::from_slice::<serde_json::Value>(&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");
}
@@ -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.