fix(scanner): make reset cleanup safely reentrant (#7180)

* chore(deps): refresh SDKs and pin clock skew regression coverage

Refresh compatible dependencies for Scanner/Heal V2 batch 1 and verify
the production S3 retry/signing path with a deterministic clock.

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

* fix(scanner): make reset cleanup safely reentrant

Refs rustfs/backlog#2264 and rustfs/backlog#2240.

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

---------

Co-authored-by: heihutu <heihutu@gmail.com>
Co-authored-by: zhi22915 <qiuzgang@gmail.com>
This commit is contained in:
houseme
2026-09-05 16:45:53 +08:00
committed by GitHub
parent e6bf2a4646
commit 42c32381b6
3 changed files with 664 additions and 55 deletions
+280 -51
View File
@@ -379,6 +379,12 @@ pub(super) fn decode_recovery_marker_for_reset(
if !matches!(marker_revision, DataUsageCacheRevision::Etag(_)) {
return Err(ScannerError::Other("cycle recovery marker has no object revision".to_string()));
}
if let Ok(value) = serde_json::from_slice::<serde_json::Value>(data)
&& let Some(state) = value.get("state")
&& !matches!(state.as_str(), Some("blocked" | "cleanup-pending"))
{
return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string()));
}
let compat = serde_json::from_slice::<ScannerCycleRecoveryMarkerCompat>(data).ok();
let _schema_version = compat.as_ref().and_then(|marker| marker.schema_version);
let primary_revision = compat
@@ -406,7 +412,10 @@ pub(super) fn decode_recovery_marker_for_reset(
};
let state = match compat.as_ref().and_then(|marker| marker.state.as_deref()) {
Some("cleanup-pending") => "cleanup-pending",
_ => "blocked",
Some("blocked") | None => "blocked",
Some(_) => {
return Err(ScannerError::Other("cycle recovery marker state is unsupported".to_string()));
}
};
let now = unix_now_secs();
Ok(ScannerCycleRecoveryMarker {
@@ -721,17 +730,19 @@ async fn mark_cycle_recovery_cleanup_pending(
mut marker: ScannerCycleRecoveryMarker,
marker_revision: &DataUsageCacheRevision,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<(ScannerCycleRecoveryMarker, DataUsageCacheRevision), ScannerError> {
marker.state = "cleanup-pending".to_string();
marker.last_attempt_at_unix_secs = unix_now_secs();
let bytes = serde_json::to_vec(&marker)
.map_err(|err| ScannerError::Other(format!("failed to encode cycle recovery marker: {err}")))?;
let info = save_config_with_publication_admission_for_epoch(
let info = save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
bytes,
marker_revision.preconditions(),
expected_epoch,
owns_reset,
)
.await
.map_err(|err| ScannerError::Other(format!("failed to mark cycle recovery cleanup pending: {err}")))?;
@@ -933,6 +944,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
.get_write_lock_quiet(Duration::from_secs(5))
.await
.map_err(|err| ScannerError::Other(format!("scanner leader lock is busy: {err}")))?;
let owns_reset = || !guard.is_lock_lost() && !ctx.is_cancelled();
if guard.is_lock_lost() {
return Err(ScannerError::Other("scanner leader lock was lost before recovery reset".to_string()));
@@ -952,7 +964,27 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
}
Err(err) => return Err(ScannerError::Other(format!("failed to read cycle recovery marker: {err}"))),
};
let marker_data = marker_data.ok_or_else(|| ScannerError::Other("scanner cycle recovery marker is absent".to_string()))?;
let Some(marker_data) = marker_data else {
// A delete may commit before its reply is lost. Confirm both durable
// fences before treating a retry without its marker as completed.
let (cycle, epoch, revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?;
let floor = persisted_usage_floor(storeapi.clone()).await?;
if !matches!(revision, DataUsageCacheRevision::Etag(_))
|| epoch < floor.leader_epoch
|| cycle.next < floor.next_cycle
|| !owns_reset()
|| scanner_publication_admission_for_epoch(storeapi.clone(), reset_epoch)
.await
.is_none()
{
return Err(ScannerError::Other(
"scanner cycle recovery marker is absent without a completed reset fence".to_string(),
));
}
set_scanner_cycle_recovery_status(recovery_status("healthy", None, false));
super::notify_scanner_cycle_recovery_wake();
return Ok(());
};
let (marker, force_full_rescan) = match serde_json::from_slice::<ScannerCycleRecoveryMarker>(&marker_data) {
Ok(marker) if validate_recovery_marker(&marker).is_ok() => (marker, false),
_ => (decode_recovery_marker_for_reset(&marker_data, &marker_revision)?, true),
@@ -1026,8 +1058,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
}
};
if let Some((primary_cycle, primary_epoch)) = primary_state {
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
let (cleanup_marker, cleanup_marker_revision) =
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch).await?;
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker.clone(), &marker_revision, reset_epoch, &owns_reset)
.await?;
set_scanner_cycle_recovery_status(recovery_status_from_marker(&cleanup_marker, "cleanup-pending"));
let usage_floor = persisted_usage_floor(storeapi.clone()).await?;
let fence_epoch = primary_epoch
@@ -1047,12 +1081,14 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"preserved scanner cycle state exceeds the bounded object size".to_string(),
));
}
let preserved_info = save_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
let preserved_info = save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
preserved_data,
primary_revision.preconditions(),
reset_epoch,
&owns_reset,
)
.await
.map_err(|err| {
@@ -1072,9 +1108,17 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"scanner leader lock was lost after fencing newer cycle state".to_string(),
));
}
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), fence_epoch, Some(reset_epoch), false)
.await
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
fence_scanner_usage_epoch_with_expected_epoch(
&ctx,
storeapi.clone(),
fence_epoch,
Some(reset_epoch),
false,
&owns_reset,
)
.await
.map_err(|err| ScannerError::Other(format!("failed to fence preserved scanner usage epoch: {err}")))?;
if guard.is_lock_lost() {
return Err(ScannerError::Other(
"scanner leader lock was lost after fencing newer cycle state".to_string(),
@@ -1088,7 +1132,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"scanner cycle state changed before recovery marker cleanup".to_string(),
));
}
delete_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &cleanup_marker_revision, &owns_reset).await?;
delete_reset_config(
storeapi.clone(),
RUSTFS_META_BUCKET,
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
@@ -1100,6 +1145,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
..Default::default()
},
reset_epoch,
&owns_reset,
)
.await
.map_err(|err| {
@@ -1149,17 +1195,20 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
// Persist the cleanup-pending phase before rewriting the primary. If the
// process dies after the rewrite, startup still sees a durable fence and
// cannot mistake the partially completed reset for a healthy state.
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
let (marker, marker_revision) = if marker.state == "cleanup-pending" {
(marker, marker_revision)
} else {
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch).await?
mark_cycle_recovery_cleanup_pending(storeapi.clone(), marker, &marker_revision, reset_epoch, &owns_reset).await?
};
let rebuilt_info = save_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
let rebuilt_info = save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
data,
primary_revision.preconditions(),
reset_epoch,
&owns_reset,
)
.await
.map_err(|err| {
@@ -1178,8 +1227,10 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
"scanner leader lock was lost after rebuilding cycle state".to_string(),
));
}
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
if let Err(err) =
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false).await
fence_scanner_usage_epoch_with_expected_epoch(&ctx, storeapi.clone(), leader_epoch, Some(reset_epoch), false, &owns_reset)
.await
{
set_scanner_cycle_recovery_status(ScannerCycleRecoveryStatus {
path: DATA_USAGE_BLOOM_NAME_PATH.clone(),
@@ -1249,7 +1300,8 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
));
}
if let Err(err) = delete_config_with_publication_admission_for_epoch(
verify_cycle_reset_intent(storeapi.clone(), &marker_revision, &owns_reset).await?;
if let Err(err) = delete_reset_config(
storeapi.clone(),
RUSTFS_META_BUCKET,
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
@@ -1261,6 +1313,7 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
..Default::default()
},
reset_epoch,
&owns_reset,
)
.await
{
@@ -1310,6 +1363,57 @@ pub async fn reset_scanner_cycle_recovery(ctx: CancellationToken, storeapi: Arc<
Ok(())
}
async fn verify_cycle_reset_intent(
storeapi: Arc<impl ScannerObjectIO>,
expected_revision: &DataUsageCacheRevision,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<(), ScannerError> {
let revision = read_config_revision(storeapi, DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to verify scanner cycle reset intent: {err}")))?;
if &revision != expected_revision {
return Err(ScannerError::Other("scanner cycle reset intent changed".to_string()));
}
if !owns_reset() {
return Err(ScannerError::Other("scanner cycle reset ownership was lost".to_string()));
}
Ok(())
}
async fn save_reset_config(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
path: &str,
data: Vec<u8>,
preconditions: crate::HTTPPreconditions,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<crate::ScannerObjectInfo, EcstoreError> {
let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else {
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
};
if !owns_reset() {
return Err(EcstoreError::other("scanner reset ownership was lost before write"));
}
save_config_with_preconditions(storeapi, path, data, preconditions).await
}
async fn delete_reset_config(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
bucket: &str,
path: &str,
options: ScannerObjectOptions,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<crate::ScannerObjectInfo, EcstoreError> {
let Some(_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), expected_epoch).await else {
return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED));
};
if !owns_reset() {
return Err(EcstoreError::other("scanner reset ownership was lost before delete"));
}
storeapi.delete_config_object(bucket, path, options).await
}
fn scanner_usage_state_reset_paths() -> Vec<String> {
vec![
DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
@@ -1333,8 +1437,14 @@ pub(super) async fn read_usage_state_reset_slots(
Ok(slots)
}
fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<PersistedUsageFloor, ScannerError> {
let mut floor = PersistedUsageFloor::default();
enum ScannerUsageResetFloor {
Missing,
Trusted(PersistedUsageFloor),
Corrupt,
}
fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<ScannerUsageResetFloor, ScannerError> {
let mut floor = None;
for slot in slots {
let Some(data) = slot.data.as_deref() else {
continue;
@@ -1342,9 +1452,21 @@ fn usage_state_reset_floor(slots: &[ScannerUsageStateResetSlot]) -> Result<Persi
let Ok(usage) = serde_json::from_slice::<DataUsageInfo>(data) else {
continue;
};
update_persisted_usage_floor(&mut floor, &usage, &slot.path)?;
if !data_usage_info_has_persisted_baseline_identity(&usage)
&& !(slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str() && data_usage_info_is_bootstrap_pending(&usage))
&& legacy_incomplete_usage_fence(data, &usage)
.and_then(|fence| fence.claimable_epoch())
.is_none()
{
continue;
}
update_persisted_usage_floor(floor.get_or_insert_with(PersistedUsageFloor::default), &usage, &slot.path)?;
}
Ok(floor)
Ok(match floor {
Some(floor) => ScannerUsageResetFloor::Trusted(floor),
None if slots.iter().any(|slot| slot.data.is_some()) => ScannerUsageResetFloor::Corrupt,
None => ScannerUsageResetFloor::Missing,
})
}
async fn read_cycle_state_for_usage_reset(
@@ -1401,11 +1523,12 @@ async fn delete_usage_state_reset_slot(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
slot: &ScannerUsageStateResetSlot,
expected_epoch: u64,
owns_reset: &(impl Fn() -> bool + Sync),
) -> Result<bool, ScannerError> {
if matches!(slot.revision, DataUsageCacheRevision::Missing) {
return Ok(false);
}
let delete_result = delete_config_with_publication_admission_for_epoch(
let delete_result = delete_reset_config(
storeapi.clone(),
RUSTFS_META_BUCKET,
&slot.path,
@@ -1415,6 +1538,7 @@ async fn delete_usage_state_reset_slot(
..Default::default()
},
expected_epoch,
owns_reset,
)
.await;
match delete_result {
@@ -1486,21 +1610,24 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary(
expected_publication_epoch: u64,
leader_epoch: Option<u64>,
context: ScannerUsageBootstrapPublishContext,
owns_publication: impl Fn() -> bool + Sync,
) -> Result<(), ScannerError> {
async fn inner(
storeapi: Arc<impl ScannerObjectIO + ScannerConfigObjectDelete>,
expected_revision: &DataUsageCacheRevision,
expected_publication_epoch: u64,
leader_epoch: Option<u64>,
owns_publication: &(impl Fn() -> bool + Sync),
) -> 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(
let save_result = save_reset_config(
storeapi.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
data.clone(),
expected_revision.preconditions(),
expected_publication_epoch,
owns_publication,
)
.await;
if save_result
@@ -1524,7 +1651,7 @@ pub(super) async fn publish_scanner_usage_bootstrap_primary(
})
}
inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch)
inner(storeapi, expected_revision, expected_publication_epoch, leader_epoch, &owns_publication)
.await
.map_err(|err| err.into_scanner_error(context))
}
@@ -1534,32 +1661,108 @@ pub(super) async fn reset_scanner_usage_state_slots_for_full_rebuild(
slots: &[ScannerUsageStateResetSlot],
expected_epoch: u64,
leader_epoch: u64,
owns_reset: impl Fn() -> bool + Sync,
) -> Result<Vec<String>, ScannerError> {
let mut reset_paths = Vec::new();
let primary = slots
.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()))?;
publish_scanner_usage_bootstrap_primary(
storeapi.clone(),
&primary.revision,
expected_epoch,
Some(leader_epoch),
ScannerUsageBootstrapPublishContext::Reset,
)
.await?;
if !owns_reset() {
return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string()));
}
let resume_epoch = usage_state_reset_resume_epoch(slots)?;
match resume_epoch {
Some(epoch) if epoch == leader_epoch => {}
Some(_) => return Err(ScannerError::Other("scanner usage reset bootstrap epoch changed".to_string())),
None => {
publish_scanner_usage_bootstrap_primary(
storeapi.clone(),
&primary.revision,
expected_epoch,
Some(leader_epoch),
ScannerUsageBootstrapPublishContext::Reset,
&owns_reset,
)
.await?;
}
}
let (data, intent_revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to inspect scanner usage reset intent: {err}")))?;
data.as_deref()
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok())
.filter(|usage| data_usage_info_is_bootstrap_pending(usage) && usage.scanner_epoch == Some(leader_epoch))
.ok_or_else(|| ScannerError::Other("scanner usage reset intent changed before cleanup".to_string()))?;
if !matches!(intent_revision, DataUsageCacheRevision::Etag(_))
|| (resume_epoch.is_some() && intent_revision != primary.revision)
{
return Err(ScannerError::Other("scanner usage reset intent revision changed".to_string()));
}
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? {
if let Some(usage) = slot
.data
.as_deref()
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok())
&& usage_epoch(&usage) >= leader_epoch
{
return Err(ScannerError::Other(format!(
"scanner usage reset slot is not older than its intent: {}",
slot.path
)));
}
let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to verify scanner usage reset intent: {err}")))?;
if revision != intent_revision {
return Err(ScannerError::Other("scanner usage reset intent changed during cleanup".to_string()));
}
if !owns_reset() {
return Err(ScannerError::Other("scanner usage reset ownership was lost".to_string()));
}
if delete_usage_state_reset_slot(storeapi.clone(), slot, expected_epoch, &owns_reset).await? {
reset_paths.push(slot.path.clone());
}
}
let revision = read_config_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.map_err(|err| ScannerError::Other(format!("failed to confirm scanner usage reset intent: {err}")))?;
if revision != intent_revision || !owns_reset() {
return Err(ScannerError::Other(
"scanner usage reset intent or ownership changed before completion".to_string(),
));
}
invalidate_admin_data_usage_snapshot_cache().await;
invalidate_data_usage_snapshot_cache().await;
Ok(reset_paths)
}
fn usage_state_reset_resume_epoch(slots: &[ScannerUsageStateResetSlot]) -> Result<Option<u64>, ScannerError> {
let primary = slots.iter().find(|slot| slot.path == DATA_USAGE_OBJ_NAME_PATH.as_str());
let usage = primary
.and_then(|slot| slot.data.as_deref())
.and_then(|data| serde_json::from_slice::<DataUsageInfo>(data).ok());
match usage {
Some(usage) if usage.usage_snapshot_bootstrap_pending => {
if !data_usage_info_is_bootstrap_pending(&usage) {
return Err(ScannerError::Other("scanner usage reset bootstrap is invalid".to_string()));
}
if usage.scanner_epoch.is_none() {
// Initial bootstrap has no reset owner yet.
return Ok(None);
}
usage
.scanner_epoch
.filter(|epoch| *epoch > 0 && *epoch < u64::MAX)
.map(Some)
.ok_or_else(|| ScannerError::Other("scanner usage reset bootstrap has no valid epoch".to_string()))
}
_ => Ok(None),
}
}
pub async fn reset_scanner_usage_state_for_full_rebuild(
ctx: CancellationToken,
storeapi: Arc<ECStore>,
@@ -1584,12 +1787,31 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
};
let (cycle, cycle_epoch, cycle_revision) = read_cycle_state_for_usage_reset(storeapi.clone()).await?;
let slots = read_usage_state_reset_slots(storeapi.clone()).await?;
let usage_floor = usage_state_reset_floor(&slots)?;
let leader_epoch = cycle_epoch
.max(usage_floor.leader_epoch)
.checked_add(1)
.filter(|epoch| *epoch < u64::MAX)
.ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?;
let usage_floor = match usage_state_reset_floor(&slots)? {
ScannerUsageResetFloor::Trusted(floor) => floor,
ScannerUsageResetFloor::Corrupt if matches!(cycle_revision, DataUsageCacheRevision::Missing) => {
return Err(ScannerError::Other("scanner usage reset has no trusted cycle or usage floor".to_string()));
}
ScannerUsageResetFloor::Missing | ScannerUsageResetFloor::Corrupt => PersistedUsageFloor {
next_cycle: cycle.next,
leader_epoch: cycle_epoch,
},
};
let resume_epoch = usage_state_reset_resume_epoch(&slots)?;
let leader_epoch = if let Some(epoch) = resume_epoch {
if epoch != cycle_epoch || usage_floor.leader_epoch > epoch || usage_floor.next_cycle > cycle.next {
return Err(ScannerError::Other(
"scanner usage reset bootstrap conflicts with the persisted cycle fence".to_string(),
));
}
epoch
} else {
cycle_epoch
.max(usage_floor.leader_epoch)
.checked_add(1)
.filter(|epoch| *epoch < u64::MAX)
.ok_or_else(|| ScannerError::Other("scanner leader epoch is exhausted".to_string()))?
};
let rebuilt_cycle = CurrentCycle {
next: cycle.next.max(usage_floor.next_cycle),
..Default::default()
@@ -1602,21 +1824,24 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
"scanner leader lock was lost before fencing usage reset cycle state".to_string(),
));
}
save_config_with_publication_admission_for_epoch(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
cycle_data,
cycle_revision.preconditions(),
reset_epoch,
)
.await
.map_err(|err| {
if scanner_publication_epoch_changed(&err) {
ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string())
} else {
ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}"))
}
})?;
if resume_epoch.is_none() {
save_reset_config(
storeapi.clone(),
DATA_USAGE_BLOOM_NAME_PATH.as_str(),
cycle_data,
cycle_revision.preconditions(),
reset_epoch,
&|| !guard.is_lock_lost() && !ctx.is_cancelled(),
)
.await
.map_err(|err| {
if scanner_publication_epoch_changed(&err) {
ScannerError::Other("scanner usage reset deferred by a movement epoch change".to_string())
} else {
ScannerError::Other(format!("failed to fence scanner cycle state for usage reset: {err}"))
}
})?;
}
if guard.is_lock_lost() {
return Err(ScannerError::Other(
@@ -1624,7 +1849,10 @@ pub async fn reset_scanner_usage_state_for_full_rebuild(
));
}
let reset_paths =
reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch).await?;
reset_scanner_usage_state_slots_for_full_rebuild(storeapi.clone(), &slots, reset_epoch, leader_epoch, || {
!guard.is_lock_lost() && !ctx.is_cancelled()
})
.await?;
if guard.is_lock_lost() {
return Err(ScannerError::Other(
"scanner leader lock was lost after publishing usage reset marker".to_string(),
@@ -2135,6 +2363,7 @@ async fn recover_legacy_incomplete_usage_floor(
expected_publication_epoch,
Some(primary.epoch),
ScannerUsageBootstrapPublishContext::Recovery,
|| true,
)
.await?;
warn!(
+7 -1
View File
@@ -191,6 +191,7 @@ pub(super) async fn initialize_usage_baseline_bootstrap(
expected_epoch,
None,
ScannerUsageBootstrapPublishContext::Initial,
|| true,
)
.await
}
@@ -201,9 +202,10 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
claimed_epoch: u64,
expected_publication_epoch: Option<u64>,
allow_bootstrap_pending: bool,
owns_fence: impl Fn() -> bool,
) -> Result<(), ScannerError> {
for retry in 0..=SCANNER_PERSIST_CAS_RETRIES {
if ctx.is_cancelled() {
if ctx.is_cancelled() || !owns_fence() {
return Err(ScannerError::Other("scanner leadership was cancelled before usage fencing".to_string()));
}
@@ -264,6 +266,9 @@ pub(super) async fn fence_scanner_usage_epoch_with_expected_epoch(
"scanner usage epoch fence changed while preparing its conditional write".to_string(),
));
};
if ctx.is_cancelled() || !owns_fence() {
return Err(ScannerError::Other("scanner leadership was lost before usage fencing".to_string()));
}
save_config_with_preconditions(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), data, revision.preconditions())
.await
};
@@ -319,6 +324,7 @@ pub(super) async fn complete_scanner_leadership_claim(
claimed_epoch,
expected_publication_epoch,
allow_bootstrap_pending,
|| true,
)
.await
{
+377 -3
View File
@@ -634,6 +634,8 @@ struct MemoryConfigStore {
cancel_after_successful_puts: Mutex<HashMap<String, (usize, CancellationToken)>>,
replace_after_successful_puts: Mutex<HashMap<String, (usize, Vec<u8>)>>,
error_after_commit_deletes: Mutex<HashSet<String>>,
cancel_after_deletes: Mutex<HashMap<String, CancellationToken>>,
pause_next_publication_admission: Mutex<Option<(Arc<tokio::sync::Notify>, Arc<tokio::sync::Notify>)>>,
put_counts: Mutex<HashMap<String, usize>>,
publication_admission_blocked: AtomicBool,
block_publication_after_admissions: AtomicUsize,
@@ -4081,6 +4083,9 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
revisions.remove(&key);
drop(revisions);
drop(objects);
if let Some(token) = self.cancel_after_deletes.lock().await.remove(&key) {
token.cancel();
}
if self.error_after_commit_deletes.lock().await.remove(&key) {
return Err(EcstoreError::other("injected delete error after commit"));
}
@@ -4088,6 +4093,11 @@ impl crate::ScannerConfigObjectDelete for MemoryConfigStore {
}
async fn scanner_data_usage_publication_admission(&self) -> Option<crate::ScannerDataUsagePublicationAdmission> {
let pause = self.pause_next_publication_admission.lock().await.take();
if let Some((entered, resume)) = pause {
entered.notify_one();
resume.notified().await;
}
if self.publication_admission_blocked.load(Ordering::Acquire) {
return None;
}
@@ -4589,7 +4599,7 @@ async fn scanner_legacy_usage_backup_survives_fencing_and_restart_after_real_met
.expect("publication must also read the intact backup");
assert_eq!(baseline.data.as_deref(), Some(data.as_slice()));
assert_eq!(baseline.revision, DataUsageCacheRevision::Missing);
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false)
fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 7, None, false, || true)
.await
.expect("legacy backup must be fenced into v2");
let fenced = read_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
@@ -4818,6 +4828,28 @@ async fn scanner_usage_state_reset_publishes_fenced_bootstrap_marker() {
);
}
let cycle_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("cycle should remain before retry");
let marker_before_retry = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("bootstrap should remain before retry");
let retry = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
.await
.expect("completed cleanup should be reentrant");
assert_eq!(retry.leader_epoch, result.leader_epoch);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("cycle should remain"),
cycle_before_retry
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("bootstrap should remain"),
marker_before_retry
);
let (floor, state) = persisted_usage_floor_for_startup(store, false)
.await
.expect("reset marker should be resumable");
@@ -4973,7 +5005,7 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() {
store.objects.lock().await.insert(key.clone(), b"newer-json".to_vec());
store.revisions.lock().await.insert(key, 2);
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3)
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true)
.await
.expect_err("stale primary revision must not be overwritten");
assert!(
@@ -4983,6 +5015,348 @@ async fn scanner_usage_state_reset_slots_reject_primary_aba() {
);
}
#[tokio::test]
async fn scanner_usage_state_reset_resumes_every_cleanup_boundary_without_rewriting_intent() {
for completed in 0..=4 {
let store = Arc::new(MemoryConfigStore::default());
let primary_path = DATA_USAGE_OBJ_NAME_PATH.as_str();
let cleanup_paths = [
format!("{primary_path}.bkp"),
LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(),
format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()),
DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str().to_string(),
];
for path in std::iter::once(primary_path).chain(cleanup_paths.iter().map(String::as_str)) {
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
usage.scanner_epoch = Some(1);
save_config(store.clone(), path, serde_json::to_vec(&usage).expect("fixture should encode"))
.await
.expect("fixture should persist");
}
// These objects belong to other owners, even when reset cleanup resumes.
for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] {
save_config(store.clone(), path, b"retain".to_vec())
.await
.expect("unrelated state should persist");
}
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
let cancelled = CancellationToken::new();
if completed == 0 {
store
.cancel_after_successful_puts
.lock()
.await
.insert(memory_config_key(RUSTFS_META_BUCKET, primary_path), (2, cancelled.clone()));
} else {
store
.cancel_after_deletes
.lock()
.await
.insert(memory_config_key(RUSTFS_META_BUCKET, &cleanup_paths[completed - 1]), cancelled.clone());
}
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled())
.await
.expect_err("interruption should stop cleanup");
assert!(err.to_string().contains("ownership"), "boundary {completed}: {err}");
for (index, path) in cleanup_paths.iter().enumerate() {
assert_eq!(
store
.objects
.lock()
.await
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, path)),
index >= completed,
"boundary {completed}, slot {index}"
);
}
let intent = read_config_with_revision(store.clone(), primary_path)
.await
.expect("intent should persist");
let slots = read_usage_state_reset_slots(store.clone())
.await
.expect("restart should reload slots");
reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
.await
.expect("restart should complete the same intent");
assert_eq!(
read_config_with_revision(store.clone(), primary_path)
.await
.expect("intent should remain"),
intent
);
assert_eq!(store.put_counts.lock().await[&memory_config_key(RUSTFS_META_BUCKET, primary_path)], 2);
for path in cleanup_paths {
assert!(
!store
.objects
.lock()
.await
.contains_key(&memory_config_key(RUSTFS_META_BUCKET, &path))
);
}
for path in ["buckets/quota-reservations/ledger", "buckets/example/incarnation"] {
assert_eq!(read_config(store.clone(), path).await.expect("unrelated state should remain"), b"retain");
}
}
}
#[tokio::test]
async fn scanner_usage_state_reset_stops_usage_fence_after_owner_loss() {
let store = Arc::new(MemoryConfigStore::default());
let mut usage = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
usage.scanner_epoch = Some(1);
let bytes = serde_json::to_vec(&usage).expect("baseline should encode");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone())
.await
.expect("baseline should persist");
let checks = AtomicUsize::new(0);
let err = fence_scanner_usage_epoch_with_expected_epoch(&CancellationToken::new(), store.clone(), 3, Some(0), false, || {
checks.fetch_add(1, Ordering::SeqCst) == 0
})
.await
.expect_err("ownership lost during reads must prevent the write");
assert!(err.to_string().contains("leadership was lost"), "{err}");
assert_eq!(
read_config(store, DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("baseline should remain"),
bytes
);
}
#[tokio::test]
async fn scanner_usage_state_reset_cancels_during_publication_admission() {
for resuming in [false, true] {
let store = Arc::new(MemoryConfigStore::default());
let usage = if resuming {
scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3))
} else {
complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0)
};
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&usage).expect("primary should encode"),
)
.await
.expect("primary should persist");
save_config(store.clone(), LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str(), b"corrupt".to_vec())
.await
.expect("cleanup target should persist");
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
let before = store.objects.lock().await.clone();
let revisions_before = store.revisions.lock().await.clone();
let entered = Arc::new(tokio::sync::Notify::new());
let resume = Arc::new(tokio::sync::Notify::new());
*store.pause_next_publication_admission.lock().await = Some((entered.clone(), resume.clone()));
let cancelled = CancellationToken::new();
let (result, ()) = tokio::join!(
reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || !cancelled.is_cancelled()),
async {
entered.notified().await;
cancelled.cancel();
resume.notify_one();
}
);
let err = result.expect_err("losing ownership during admission must prevent mutation");
assert!(err.to_string().contains("ownership was lost"), "resuming={resuming}: {err}");
assert_eq!(*store.objects.lock().await, before);
assert_eq!(*store.revisions.lock().await, revisions_before);
}
}
#[tokio::test]
#[serial]
async fn scanner_usage_state_reset_rejects_corruption_without_a_trusted_floor() {
let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), b"{corrupt".to_vec())
.await
.expect("corrupt primary should persist");
let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("evidence should load");
let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
.await
.expect_err("corruption must not become a zero floor");
assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("evidence should remain"),
before
);
assert!(matches!(
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
let mut backup = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
backup.scanner_epoch = Some(7);
backup.scanner_cycle = Some(40);
save_config(
store.clone(),
&format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()),
serde_json::to_vec(&backup).expect("backup should encode"),
)
.await
.expect("valid backup should persist");
let result = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store)
.await
.expect("valid backup should supply the recovery floor");
assert_eq!(result.leader_epoch, 8);
assert_eq!(result.next_cycle, 41);
}
#[tokio::test]
async fn scanner_usage_state_reset_rejects_replaced_intent_and_newer_cleanup_slot() {
let store = Arc::new(MemoryConfigStore::default());
let marker = scanner_usage_bootstrap_marker(std::time::SystemTime::UNIX_EPOCH, Some(3));
let bytes = serde_json::to_vec(&marker).expect("marker should encode");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes.clone())
.await
.expect("intent should persist");
let slots = read_usage_state_reset_slots(store.clone()).await.expect("slots should load");
save_config(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str(), bytes)
.await
.expect("another intent should persist");
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
.await
.expect_err("same epoch cannot replace an intent revision");
assert!(err.to_string().contains("intent revision changed"), "{err}");
let mut newer = complete_usage_with_bucket_count(Some(std::time::SystemTime::UNIX_EPOCH), 0);
newer.scanner_epoch = Some(3);
let path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str());
let bytes = serde_json::to_vec(&newer).expect("newer snapshot should encode");
save_config(store.clone(), &path, bytes.clone())
.await
.expect("newer snapshot should persist");
let slots = read_usage_state_reset_slots(store.clone())
.await
.expect("slots should reload");
let err = reset_scanner_usage_state_slots_for_full_rebuild(store.clone(), &slots, 0, 3, || true)
.await
.expect_err("cleanup cannot delete same-epoch progress");
assert!(err.to_string().contains("not older than its intent"), "{err}");
assert_eq!(read_config(store, &path).await.expect("newer snapshot should remain"), bytes);
}
#[tokio::test]
#[serial]
async fn scanner_usage_state_reset_rejects_decodable_untrusted_floor() {
let (_temp_dir, store) = setup_scanner_cycle_store_with_usage_baseline(false).await;
let invalid_identity = DataUsageInfo {
usage_snapshot_complete: true,
buckets_count: 1,
last_update: Some(std::time::SystemTime::UNIX_EPOCH),
..Default::default()
};
for usage in [DataUsageInfo::default(), invalid_identity] {
save_config(
store.clone(),
DATA_USAGE_OBJ_NAME_PATH.as_str(),
serde_json::to_vec(&usage).expect("fixture should encode"),
)
.await
.expect("untrusted primary should persist");
let before = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("primary should load");
let err = reset_scanner_usage_state_for_full_rebuild(CancellationToken::new(), store.clone())
.await
.expect_err("valid JSON alone cannot prove a usage floor");
assert!(err.to_string().contains("no trusted cycle or usage floor"), "{err}");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("evidence should remain"),
before
);
assert!(matches!(
read_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str()).await,
Err(EcstoreError::ConfigNotFound)
));
}
}
#[test]
fn full_rescan_reset_rejects_unknown_marker_phase_even_with_invalid_compat_fields() {
for state in [serde_json::json!("rewrite-v2"), serde_json::json!(7), serde_json::Value::Null] {
let marker = serde_json::json!({"state": state, "retry_count": "future-type", "schema_version": 99});
let err = super::cycle_state::decode_recovery_marker_for_reset(
&serde_json::to_vec(&marker).expect("future marker should encode"),
&DataUsageCacheRevision::Etag("intent-1".to_string()),
)
.expect_err("unknown persistent phases must remain fenced");
assert!(err.to_string().contains("state is unsupported"), "{err}");
}
}
#[tokio::test]
#[serial]
async fn full_rescan_reset_preserves_unknown_phase_and_retries_completed_cleanup() {
let (_temp_dir, store) = setup_scanner_cycle_store().await;
save_config(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str(), b"corrupt".to_vec())
.await
.expect("corrupt primary should persist");
save_config(
store.clone(),
DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(),
br#"{"state":"future-rewrite"}"#.to_vec(),
)
.await
.expect("future marker should persist");
let primary_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary should load");
let marker_before = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("marker should load");
let err = reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect_err("unknown phase must block explicit reset");
assert!(err.to_string().contains("state is unsupported"), "{err}");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("primary should remain"),
primary_before
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str())
.await
.expect("marker should remain"),
marker_before
);
save_config(store.clone(), DATA_USAGE_BLOOM_RECOVERY_PATH.as_str(), b"{malformed".to_vec())
.await
.expect("recoverable marker should persist");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("reset should complete");
let primary = read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt primary should load");
let usage = read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("fenced usage should load");
reset_scanner_cycle_recovery(CancellationToken::new(), store.clone())
.await
.expect("retry after marker deletion should complete");
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_BLOOM_NAME_PATH.as_str())
.await
.expect("rebuilt primary should remain"),
primary
);
assert_eq!(
read_config_with_revision(store.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str())
.await
.expect("fenced usage should remain"),
usage
);
}
#[tokio::test]
async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() {
let store = Arc::new(MemoryConfigStore::default());
@@ -4993,7 +5367,7 @@ async fn scanner_usage_state_reset_slots_defer_when_publication_epoch_moves() {
.expect("usage reset slots should be inspected");
store.publication_admission_blocked.store(true, Ordering::Release);
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3)
let err = reset_scanner_usage_state_slots_for_full_rebuild(store, &slots, 0, 3, || true)
.await
.expect_err("movement admission loss must defer reset");
assert!(