From 2dcc064a5a4d8612a30621b1bf3cf2cb9397ff05 Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 30 Aug 2026 03:00:24 +0800 Subject: [PATCH] fix(scanner): fence backup and cleanup mutations Carry storage-owned publication scopes through backup CAS writes and observed cleanup deletes so cancellation and absolute lease deadlines retain the same safety boundary. Co-Authored-By: heihutu --- 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, 158 insertions(+), 172 deletions(-) diff --git a/crates/ecstore/src/set_disk/ops/object.rs b/crates/ecstore/src/set_disk/ops/object.rs index b4e349b39..5565a3f51 100644 --- a/crates/ecstore/src/set_disk/ops/object.rs +++ b/crates/ecstore/src/set_disk/ops/object.rs @@ -270,6 +270,22 @@ 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); @@ -7098,6 +7114,11 @@ 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. @@ -7192,6 +7213,7 @@ 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; @@ -7206,6 +7228,7 @@ 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 { @@ -7217,10 +7240,14 @@ 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()); } @@ -7228,10 +7255,14 @@ 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()); } @@ -7304,10 +7335,14 @@ 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)); } @@ -7381,6 +7416,7 @@ 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]))?; @@ -7392,6 +7428,9 @@ 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); } @@ -7417,6 +7456,7 @@ 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]))?; @@ -7442,6 +7482,9 @@ 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 9091bd6d9..3a8cb8d6f 100644 --- a/crates/scanner/src/lib.rs +++ b/crates/scanner/src/lib.rs @@ -753,6 +753,32 @@ 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, @@ -763,10 +789,7 @@ pub(crate) async fn delete_config_with_publication_admission_for_epoch( where S: ScannerObjectIO + ScannerConfigObjectDelete, { - 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 + delete_config_with_publication_scope_for_epoch(api, bucket, object, opts, expected_epoch, None).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 f213165e3..1139b5345 100644 --- a/crates/scanner/src/scanner.rs +++ b/crates/scanner/src/scanner.rs @@ -37,6 +37,7 @@ 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; @@ -455,7 +456,6 @@ 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 } -#[allow(dead_code)] +#[cfg(test)] 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_and_scope( + sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_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_and_scope( +async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope( ctx: &CancellationToken, storeapi: Arc, expected_publication_epoch: Option, remote_lease_deadline: Option, scanner_publication_lease_fence: Option<&str>, - remote_lease_tokens: Vec, + remote_lease_tokens: &[Uuid], lease_release_safe: Arc, ) -> Result<(), EcstoreError> { let backup_path = format!("{}.bkp", DATA_USAGE_OBJ_NAME_PATH.as_str()); @@ -549,31 +549,32 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_s if remote_lease_deadline.is_some_and(|deadline| std::time::Instant::now() >= deadline) { return Err(EcstoreError::other(SCANNER_PUBLICATION_EPOCH_CHANGED)); } - 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 = 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 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, @@ -584,19 +585,20 @@ async fn sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_and_s publication_scope.clone(), ) .await; + drop(legacy_publication_admission); if let Some(scope) = publication_scope { match scope.wait_for_completion().await { - 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", - )), + 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", + )); + } } - } else { - save_result } + save_result }; match save_result { @@ -1473,10 +1475,13 @@ 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_data_usage_persist_baseline(storeapi.clone()).await; + let usage_persist_baseline_result = read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await; drop(baseline_publication_guard); let usage_persist_baseline = match usage_persist_baseline_result { - Ok(baseline) => baseline, + Ok((data, revision)) => DataUsagePersistBaseline { + data: data.map(Bytes::from), + revision, + }, Err(err) => { error!( target: "rustfs::scanner", @@ -1516,20 +1521,6 @@ 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), }; @@ -2879,7 +2870,12 @@ 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_for_result(&scan_cycle_result, usage_persist_outcome); + 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 pending_maintenance_work = scan_cycle_result.has_pending_maintenance_work(); let durable_complete_snapshot = scan_cycle_result.status == ScannerCycleStatus::Complete && matches!( @@ -2894,34 +2890,6 @@ 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 674f43b3f..617c2f016 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) } -pub(super) fn scanner_publication_scope_deadline( +fn scanner_publication_scope_deadline( persist_timeout: Duration, remote_lease_deadline: Option, ) -> tokio::time::Instant { @@ -53,84 +53,6 @@ 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. @@ -388,8 +310,8 @@ where publication_epoch = Some(read_epoch); let authoritative_data = match next_baseline.as_ref() { Some(baseline) => baseline.data.clone(), - None => match read_data_usage_persist_baseline(storeapi.clone()).await { - Ok(baseline) => baseline.data, + None => match read_config_with_revision(storeapi.clone(), DATA_USAGE_OBJ_NAME_PATH.as_str()).await { + Ok((data, _)) => data.map(Bytes::from), Err(err) => { error!( target: "rustfs::scanner", @@ -754,6 +676,8 @@ 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 { @@ -777,6 +701,8 @@ 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 { @@ -819,6 +745,8 @@ 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 { @@ -836,13 +764,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_and_scope( + let backup_result = sync_data_usage_backup_from_primary_for_epoch_and_lease_and_fence_with_scope( &ctx, storeapi.clone(), expected_publication_epoch, remote_lease_deadline, scanner_publication_lease_fence.as_deref(), - remote_lease_tokens.clone(), + &remote_lease_tokens, Arc::clone(&lease_release_safe), ) .await; @@ -877,6 +805,8 @@ 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; @@ -945,7 +875,15 @@ async fn cleanup_observed_data_usage_snapshot_for_epoch_and_lease( return false; } - let result = delete_config_with_publication_admission_for_epoch( + 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( storeapi, RUSTFS_META_BUCKET, DATA_USAGE_OBSERVED_OBJ_NAME_PATH.as_str(), @@ -964,9 +902,23 @@ 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(