From 290deb866510ce26fd415050104d5d96e2966958 Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 30 Aug 2026 03:02:03 +0800 Subject: [PATCH] Revert "fix(scanner): fence backup and cleanup mutations" This reverts commit 2dcc064a5a4d8612a30621b1bf3cf2cb9397ff05. --- crates/ecstore/src/set_disk/ops/object.rs | 43 ------- crates/scanner/src/lib.rs | 31 +---- crates/scanner/src/scanner.rs | 136 +++++++++++++--------- crates/scanner/src/scanner/usage_store.rs | 120 +++++++++++++------ 4 files changed, 172 insertions(+), 158 deletions(-) diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index 5565a3f51..b4e349b39 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -270,22 +270,6 @@ const OLD_DATA_CLEANUP_RECEIPT_FILE: &str = ".rustfs-old-data-cleanup-receipt.js const SCANNER_PUBLICATION_LEASE_FENCE_MAX_BYTES: usize = 64 * 1024; const SCANNER_PUBLICATION_LEASE_FENCE_MAX_ENTRIES: usize = 256; -fn begin_scanner_publication_delete_mutation(scope: Option<&crate::object_api::ScannerPublicationCommitScope>) -> Result<()> { - let Some(scope) = scope else { - return Ok(()); - }; - if scope.state() == crate::object_api::ScannerPublicationCommitState::Admitted { - scope - .try_begin() - .map_err(|err| Error::other(format!("scanner publication delete scope cannot start: {err:?}")))?; - } - if !scope.can_commit() { - let _ = scope.mark_indeterminate(); - return Err(StorageError::OperationCanceled); - } - Ok(()) -} - fn take_scanner_publication_lease_tokens(user_defined: &mut HashMap) -> Result>> { let Some(encoded) = user_defined.remove(SCANNER_PUBLICATION_LEASE_FENCE_METADATA_KEY) else { return Ok(None); @@ -7114,11 +7098,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { #[tracing::instrument(skip(self, opts))] async fn delete_object(&self, bucket: &str, object: &str, mut opts: ObjectOptions) -> Result { - let _scope_outcome_guard = opts - .scanner_publication_commit_scope - .clone() - .map(ScannerPublicationCommitScopeGuard::new); - let scanner_publication_commit_scope = opts.scanner_publication_commit_scope.clone(); // Scanner cleanup carries the per-peer lease fence as transient // request metadata. Consume it before any delete-prefix fanout so it // cannot be persisted or treated as user metadata. @@ -7213,7 +7192,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } delete_request.set_skip_tier_free_version(); } - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_object_version(bucket, object, &delete_request, false).await?; if let Some((_, deleted_object)) = replication_delete { ReplicationLifecycleBridge::schedule_delete(bucket.to_string(), deleted_object).await; @@ -7228,7 +7206,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { ..Default::default() }; delete_request.set_tier_free_version_id(&Uuid::new_v4().to_string()); - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_object_version(bucket, object, &delete_request, false).await?; } for version in &versions.free_versions { @@ -7240,14 +7217,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { ..Default::default() }; delete_request.set_tier_free_version(); - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_object_version(bucket, object, &delete_request, false).await?; } } } - if let Some(scope) = scanner_publication_commit_scope.as_ref() { - let _ = scope.mark_committed(); - } self.invalidate_get_object_metadata_cache(bucket, object).await; return Ok(ObjectInfo::default()); } @@ -7255,14 +7228,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { self.validate_bucket_incarnation(bucket, expected_incarnation_id).await?; } ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_prefix_with_scanner_publication_lease(bucket, object, scanner_publication_lease_tokens.as_ref()) .await .map_err(|e| to_object_err(e.into(), vec![bucket, object]))?; - if let Some(scope) = scanner_publication_commit_scope.as_ref() { - let _ = scope.mark_committed(); - } self.invalidate_all_get_object_metadata_cache(); return Ok(ObjectInfo::default()); } @@ -7335,14 +7304,10 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { ..Default::default() }; ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_object_version(bucket, object, &dfi, false) .await .map_err(|e| to_object_err(e, vec![bucket, object]))?; self.invalidate_get_object_metadata_cache(bucket, object).await; - if let Some(scope) = scanner_publication_commit_scope.as_ref() { - let _ = scope.mark_committed(); - } return Ok(ObjectInfo::from_file_info(&dfi, bucket, object, opts.versioned || opts.version_suspended)); } @@ -7416,7 +7381,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { }; ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_object_version(bucket, object, &fi, should_force_delete_marker_for_missing_version(&opts)) .await .map_err(|e| to_object_err(e, vec![bucket, object]))?; @@ -7428,9 +7392,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { oi.user_tags = Arc::clone(&goi.user_tags); oi.replication_decision = goi.replication_decision; self.invalidate_get_object_metadata_cache(bucket, object).await; - if let Some(scope) = scanner_publication_commit_scope.as_ref() { - let _ = scope.mark_committed(); - } return Ok(oi); } @@ -7456,7 +7417,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { } ensure_delete_commit_locks_held(_lock_guard.as_ref(), bucket, object, &opts)?; - begin_scanner_publication_delete_mutation(scanner_publication_commit_scope.as_ref())?; self.delete_object_version(bucket, object, &dfi, opts.delete_marker) .await .map_err(|e| to_object_err(e, vec![bucket, object]))?; @@ -7482,9 +7442,6 @@ impl crate::storage_api_contracts::object::ObjectOperations for SetDisks { obj_info.delete_marker = true; } self.invalidate_get_object_metadata_cache(bucket, object).await; - if let Some(scope) = scanner_publication_commit_scope.as_ref() { - let _ = scope.mark_committed(); - } Ok(obj_info) } diff --git a/crates/scanner/src/lib.rs b/crates/scanner/src/lib.rs index 3a8cb8d6f..9091bd6d9 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -753,32 +753,6 @@ pub(crate) fn scanner_publication_epoch_changed(error: &EcstoreError) -> bool { ) } -pub(crate) async fn delete_config_with_publication_scope_for_epoch( - api: Arc, - bucket: &str, - object: &str, - mut opts: ScannerObjectOptions, - expected_epoch: u64, - scanner_publication_commit_scope: Option, -) -> EcstoreResult -where - S: ScannerObjectIO + ScannerConfigObjectDelete, -{ - let legacy_admission = if scanner_publication_commit_scope.is_none() { - Some( - scanner_publication_admission_for_epoch(api.clone(), expected_epoch) - .await - .ok_or_else(|| EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED))?, - ) - } else { - None - }; - opts.scanner_publication_commit_scope = scanner_publication_commit_scope; - let result = api.delete_config_object(bucket, object, opts).await; - drop(legacy_admission); - result -} - pub(crate) async fn delete_config_with_publication_admission_for_epoch( api: Arc, bucket: &str, @@ -789,7 +763,10 @@ pub(crate) async fn delete_config_with_publication_admission_for_epoch( where S: ScannerObjectIO + ScannerConfigObjectDelete, { - delete_config_with_publication_scope_for_epoch(api, bucket, object, opts, expected_epoch, None).await + let Some(_admission) = scanner_publication_admission_for_epoch(api.clone(), expected_epoch).await else { + return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); + }; + api.delete_config_object(bucket, object, opts).await } /// Capture the storage-owned publication epoch without retaining the read diff --git a/crates/scanner/src/scanner.rs b/crates/scanner/src/scanner.rs index 1139b5345..f213165e3 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -37,7 +37,6 @@ use crate::scanner_io::{ scanner_maintenance_generation, }; use crate::sleeper::{SCANNER_SLEEPER, set_scanner_default_speed}; -use crate::storage_api::owner::ScannerPublicationCommitState; use crate::{DataUsageInfo, ScannerActivityGuard, ScannerError, ScannerRuntimeGuard}; use crate::{ScannerConfigObjectDelete, ScannerObjectIO, ScannerObjectOptions}; use bytes::Bytes; @@ -456,6 +455,7 @@ fn data_usage_backup_due(data_usage_info: &DataUsageInfo) -> bool { } #[cfg(test)] +#[allow(dead_code)] async fn sync_data_usage_backup_from_primary( ctx: &CancellationToken, storeapi: Arc, @@ -463,7 +463,7 @@ async fn sync_data_usage_backup_from_primary( sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence(ctx, storeapi, None, None, None).await } -#[cfg(test)] +#[allow(dead_code)] async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence( ctx: &CancellationToken, storeapi: Arc, @@ -471,25 +471,25 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence( remote_lease_deadline: Option, scanner_publication_lease_fence: Option<&str>, ) -> Result<(), EcstoreError> { - sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope( + sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_scope( ctx, storeapi, expected_publication_epoch, remote_lease_deadline, scanner_publication_lease_fence, - &[], + Vec::new(), Arc::new(AtomicBool::new(true)), ) .await } -async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope( +async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_scope( ctx: &CancellationToken, storeapi: Arc, expected_publication_epoch: Option, remote_lease_deadline: Option, scanner_publication_lease_fence: Option<&str>, - remote_lease_tokens: &[Uuid], + remote_lease_tokens: Vec, lease_release_safe: Arc, ) -> Result<(), EcstoreError> { let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -549,32 +549,31 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_ if remote_lease_deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) { return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); } - let publication_scope = storeapi - .scanner_data_usage_publication_commit_scope_with_release_flag( - read_epoch, - tokio::time::Instant::now() - .checked_add(data_usage_persist_timeout()) - .unwrap_or_else(tokio::time::Instant::now) - .min( - remote_lease_deadline - .map(tokio::time::Instant::from_std) - .unwrap_or_else(|| tokio::time::Instant::now() + data_usage_persist_timeout()), - ), - remote_lease_tokens.to_vec(), - Arc::clone(&lease_release_safe), - ) - .await; - let legacy_publication_admission = if publication_scope.is_none() { - let Some(admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { - if retry < SCANNER_PERSIST_CAS_RETRIES { - continue; - } - return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); - }; - Some(admission) - } else { - None + let Some(_publication_admission) = scanner_publication_admission_for_epoch(storeapi.clone(), read_epoch).await else { + if retry < SCANNER_PERSIST_CAS_RETRIES { + continue; + } + return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); }; + let publication_scope = match expected_publication_epoch { + Some(expected_epoch) => { + storeapi + .scanner_data_usage_publication_commit_scope_with_release_flag( + expected_epoch, + usage_store::scanner_publication_scope_deadline(data_usage_persist_timeout(), remote_lease_deadline), + remote_lease_tokens.clone(), + Arc::clone(&lease_release_safe), + ) + .await + } + None => None, + }; + if expected_publication_epoch.is_some() && publication_scope.is_none() { + if retry < SCANNER_PERSIST_CAS_RETRIES { + continue; + } + return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); + } let save_result = save_config_shared_with_preconditions_and_lease_fence_and_scope( storeapi.clone(), &backup_path, @@ -585,20 +584,19 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_ publication_scope.clone(), ) .await; - drop(legacy_publication_admission); if let Some(scope) = publication_scope { match scope.wait_for_completion().await { - ScannerPublicationCommitState::Committed | ScannerPublicationCommitState::AbortedBeforeCommit => {} - ScannerPublicationCommitState::Indeterminate - | ScannerPublicationCommitState::Admitted - | ScannerPublicationCommitState::InFlight => { - return Err(EcstoreError::other( - "scanner backup publication scope did not reach a safe terminal state", - )); - } + crate::storage_api::owner::ScannerPublicationCommitState::Committed + | crate::storage_api::owner::ScannerPublicationCommitState::AbortedBeforeCommit => save_result, + crate::storage_api::owner::ScannerPublicationCommitState::Indeterminate + | crate::storage_api::owner::ScannerPublicationCommitState::Admitted + | crate::storage_api::owner::ScannerPublicationCommitState::InFlight => Err(EcstoreError::other( + "scanner backup publication commit scope did not reach a safe terminal state", + )), } + } else { + save_result } - save_result }; match save_result { @@ -1475,13 +1473,10 @@ async fn run_data_scanner_cycle_with_budget( mark_scan_cycle_idle(cycle_info, &mut cycle_metrics_guard).await; return ScannerCycleOutcome::Deferred(ScannerCycleDeferReason::DataMovement); }; - let usage_persist_baseline_result = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await; + let usage_persist_baseline_result = read_data_usage_persist_baseline(storeapi.clone()).await; drop(baseline_publication_guard); let usage_persist_baseline = match usage_persist_baseline_result { - Ok((data, revision)) => DataUsagePersistBaseline { - data: data.map(Bytes::from), - revision, - }, + Ok(baseline) => baseline, Err(err) => { error!( target: "rustfs::scanner", @@ -1521,6 +1516,20 @@ async fn run_data_scanner_cycle_with_budget( { Some(ScannerCycleDeferReason::DataMovement) } + // A complete walk can still be retained as an observational snapshot + // when only the final activity proof was unavailable. It must not + // block the observation receiver: the authoritative publication + // fence remains enforced by the usage store and the cycle is advanced + // as partial without acknowledging dirty usage. + Ok(result) + if result.has_observational_snapshot() + && matches!( + result.status, + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + ) => + { + None + } Ok(result) => final_data_usage_publication_defer_reason(storeapi.as_ref(), result.status).await, Err(_) => Some(ScannerCycleDeferReason::ActivityBaselineUnavailable), }; @@ -2870,12 +2879,7 @@ fn finalize_scanner_cycle_result( scan_cycle_result: crate::scanner_io::ScannerCycleResult, usage_persist_outcome: DataUsagePersistOutcome, ) -> (ScannerCycleOutcome, bool, Vec) { - let completion_outcome = scanner_cycle_completion_outcome( - scan_cycle_result.status, - usage_persist_outcome, - scan_cycle_result.has_dirty_usage_to_acknowledge(), - scan_cycle_result.has_failed_dirty_usage(), - ); + let completion_outcome = scanner_cycle_completion_outcome_for_result(&scan_cycle_result, usage_persist_outcome); let pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work(); let durable_complete_snapshot = scan_cycle_result.status == ScannerCycleStatus::Complete && matches!( @@ -2890,6 +2894,34 @@ fn finalize_scanner_cycle_result( (completion_outcome, pending_maintenance_work, remote_dirty_usage_acknowledgements) } +fn scanner_cycle_completion_outcome_for_result( + scan_cycle_result: &crate::scanner_io::ScannerCycleResult, + usage_persist_outcome: DataUsagePersistOutcome, +) -> ScannerCycleOutcome { + let has_dirty_usage = scan_cycle_result.has_dirty_usage_to_acknowledge(); + let has_failed_dirty_usage = scan_cycle_result.has_failed_dirty_usage(); + if scan_cycle_result.has_observational_snapshot() + && matches!( + scan_cycle_result.status, + ScannerCycleStatus::Deferred(ScannerCycleDeferReason::ActivityBaselineUnavailable) + ) + { + return match usage_persist_outcome { + DataUsagePersistOutcome::Saved + | DataUsagePersistOutcome::AlreadyDurable + | DataUsagePersistOutcome::PriorCycleDurable + | DataUsagePersistOutcome::Current + if !has_failed_dirty_usage => + { + ScannerCycleOutcome::Partial + } + DataUsagePersistOutcome::Deferred(reason) => ScannerCycleOutcome::Deferred(reason), + _ => ScannerCycleOutcome::Failed, + }; + } + scanner_cycle_completion_outcome(scan_cycle_result.status, usage_persist_outcome, has_dirty_usage, has_failed_dirty_usage) +} + /// Decide whether an incoming usage snapshot must be skipped as stale, given the local /// wall clock `now`. Mirrors `stale_data_usage_persist_reason` in /// `crates/ecstore/src/data_usage/mod.rs` — keep the two consistent. diff --git a/crates/scanner/src/scanner/usage_store.rs b/crates/scanner/src/scanner/usage_store.rs index 617c2f016..674f43b3f 100644 --- a/crates/scanner/src/scanner/usage_store.rs +++ b/crates/scanner/src/scanner/usage_store.rs @@ -37,7 +37,7 @@ fn remote_lease_expired(deadline: Option) -> bool { deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) } -fn scanner_publication_scope_deadline( +pub(super) fn scanner_publication_scope_deadline( persist_timeout: Duration, remote_lease_deadline: Option, ) -> tokio::time::Instant { @@ -53,6 +53,84 @@ pub(super) struct DataUsagePersistBaseline { pub(super) revision: DataUsageCacheRevision, } +/// Read the bytes used as the baseline for a usage publication while keeping +/// the v2 primary revision as the CAS fence. During an interrupted upgrade the +/// primary can be valid JSON without a baseline identity; in that case a +/// same-or-newer durable companion may still be used, but an older legacy +/// snapshot must not cross the primary's epoch fence. +pub(super) async fn read_data_usage_persist_baseline( + storeapi: Arc, +) -> Result { + let (primary, revision) = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await?; + let Some(primary) = primary else { + for path in [ + format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()), + LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(), + format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()), + ] { + let (candidate, _) = read_config_with_revision(storeapi.clone(), &path).await?; + let Some(candidate) = candidate else { + continue; + }; + let Ok(usage) = serde_json::from_slice::(&candidate) else { + continue; + }; + if data_usage_info_has_persisted_baseline_identity(&usage) { + return Ok(DataUsagePersistBaseline { + data: Some(Bytes::from(candidate)), + revision, + }); + } + } + return Ok(DataUsagePersistBaseline { data: None, revision }); + }; + + let Ok(primary_info) = serde_json::from_slice::(&primary) else { + // Preserve the original bytes and revision. A completed scan may + // replace the invalid primary under this CAS fence; an observation + // will still reject it below because it has no verifiable identity. + return Ok(DataUsagePersistBaseline { + data: Some(Bytes::from(primary)), + revision, + }); + }; + if data_usage_info_has_persisted_baseline_identity(&primary_info) || data_usage_info_is_bootstrap_pending(&primary_info) { + return Ok(DataUsagePersistBaseline { + data: Some(Bytes::from(primary)), + revision, + }); + } + + let invalid_primary_epoch = primary_info.scanner_epoch; + for path in [ + format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()), + LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str().to_string(), + format!("{}.bkp", LEGACY_DATA_USAGE_OBJ_NAME_PATH.as_str()), + ] { + let (candidate, _) = read_config_with_revision(storeapi.clone(), &path).await?; + let Some(candidate) = candidate else { + continue; + }; + let Ok(usage) = serde_json::from_slice::(&candidate) else { + continue; + }; + let candidate_epoch = usage.scanner_epoch.unwrap_or_default(); + if data_usage_info_has_persisted_baseline_identity(&usage) + && invalid_primary_epoch.is_none_or(|epoch| candidate_epoch >= epoch) + { + return Ok(DataUsagePersistBaseline { + data: Some(Bytes::from(candidate)), + revision, + }); + } + } + + Ok(DataUsagePersistBaseline { + data: Some(Bytes::from(primary)), + revision, + }) +} + /// Short-lived publication inputs captured for one usage persistence attempt. /// Keeping the movement epoch, lease deadline, and target fence together makes /// it explicit that they are one proof rather than independent options. @@ -310,8 +388,8 @@ where publication_epoch = Some(read_epoch); let authoritative_data = match next_baseline.as_ref() { Some(baseline) => baseline.data.clone(), - None => match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await { - Ok((data, _)) => data.map(Bytes::from), + None => match read_data_usage_persist_baseline(storeapi.clone()).await { + Ok(baseline) => baseline.data, Err(err) => { error!( target: "rustfs::scanner", @@ -676,8 +754,6 @@ where expected_publication_epoch, remote_lease_deadline, scanner_publication_lease_fence.as_deref(), - &remote_lease_tokens, - Arc::clone(&lease_release_safe), ) .await; if expected_publication_epoch.is_some() && !cleanup_ok { @@ -701,8 +777,6 @@ where expected_publication_epoch, remote_lease_deadline, scanner_publication_lease_fence.as_deref(), - &remote_lease_tokens, - Arc::clone(&lease_release_safe), ) .await; if expected_publication_epoch.is_some() && !cleanup_ok { @@ -745,8 +819,6 @@ where expected_publication_epoch, remote_lease_deadline, scanner_publication_lease_fence.as_deref(), - &remote_lease_tokens, - Arc::clone(&lease_release_safe), ) .await; if expected_publication_epoch.is_some() && !cleanup_ok { @@ -764,13 +836,13 @@ where if backup_due { let done_save = Metrics::time(Metric::SaveUsage); - let backup_result = sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope( + let backup_result = sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_scope( &ctx, storeapi.clone(), expected_publication_epoch, remote_lease_deadline, scanner_publication_lease_fence.as_deref(), - &remote_lease_tokens, + remote_lease_tokens.clone(), Arc::clone(&lease_release_safe), ) .await; @@ -805,8 +877,6 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease( expected_publication_epoch: Option, remote_lease_deadline: Option, scanner_publication_lease_fence: Option<&str>, - remote_lease_tokens: &[Uuid], - lease_release_safe: Arc, ) -> bool { if remote_lease_expired(remote_lease_deadline) { return false; @@ -875,15 +945,7 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease( return false; } - let publication_scope = storeapi - .scanner_data_usage_publication_commit_scope_with_release_flag( - read_epoch, - scanner_publication_scope_deadline(data_usage_persist_timeout(), remote_lease_deadline), - remote_lease_tokens.to_vec(), - Arc::clone(&lease_release_safe), - ) - .await; - let result = crate::delete_config_with_publication_scope_for_epoch( + let result = delete_config_with_publication_admission_for_epoch( storeapi, RUSTFS_META_BUCKET, DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str(), @@ -902,23 +964,9 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease( ..Default::default() }, read_epoch, - publication_scope.clone(), ) .await; - let result = if let Some(scope) = publication_scope { - match scope.wait_for_completion().await { - ScannerPublicationCommitState::Committed | ScannerPublicationCommitState::AbortedBeforeCommit => result, - ScannerPublicationCommitState::Indeterminate - | ScannerPublicationCommitState::Admitted - | ScannerPublicationCommitState::InFlight => Err(EcstoreError::other( - "scanner publication cleanup scope did not reach a safe terminal state", - )), - } - } else { - result - }; - match result { Ok(_) | Err(