From fab2d7e144118ad6cb53f7ad882aa298a635346b Mon Sep 17 00:00:00 2001 From: Chris Date: Wed, 30 Sep 2026 15:45:37 +0800 Subject: [PATCH] fix(scanner): serialize pause backlog replica publication (#8264) --- crates/scanner/src/scanner/backlog.rs | 70 +++++++++++++++++++++------ crates/scanner/src/storage_api.rs | 6 ++- 2 files changed, 60 insertions(+), 16 deletions(-) diff --git a/crates/scanner/src/scanner/backlog.rs b/crates/scanner/src/scanner/backlog.rs index 9229e69f6..b8f2eb314 100644 --- a/crates/scanner/src/scanner/backlog.rs +++ b/crates/scanner/src/scanner/backlog.rs @@ -1587,26 +1587,24 @@ where .iter() .filter_map(|replica| replica.revision.clone().map(|revision| (replica.id, revision))) .collect::>(); - let results = join_all(writable.into_iter().map(|set| { + // Replica writes can share the pool namespace or the fixed multipool lock. + // Finish each write before starting another acquisition for this publication. + let mut results = Vec::with_capacity(writable.len()); + for set in writable { let id = ScannerPauseBacklogReplicaId { pool_index: set.pool_index, set_index: set.set_index, }; - let revision = revisions.get(&id).cloned(); - let data = data.clone(); - let storeapi = storeapi.clone(); - async move { - let Some(revision) = revision else { - return (id, Err("replica revision is unavailable".to_string())); - }; - let result = storeapi - .save_scanner_pause_backlog_replica(id.pool_index, id.set_index, data, revision.preconditions()) + let result = match revisions.get(&id) { + Some(revision) => storeapi + .clone() + .save_scanner_pause_backlog_replica(id.pool_index, id.set_index, data.clone(), revision.preconditions()) .await - .map_err(|err| err.to_string()); - (id, result) - } - })) - .await; + .map_err(|err| err.to_string()), + None => Err("replica revision is unavailable".to_string()), + }; + results.push((id, result)); + } let failures = results .iter() @@ -2243,6 +2241,48 @@ mod tests { } } + #[test] + #[serial_test::serial] + fn native_writer_does_not_spend_replica_lock_budget_waiting_for_its_own_write() { + run_native_retirement_test(async || { + use crate::storage_api::owner::{PutObjectCommitBarrier, PutObjectCommitPause, ecstore_get_lock_acquire_timeout}; + + let (_root, store) = native_retirement_store().await; + let before = load_scanner_pause_backlog(Arc::clone(&store)) + .await + .expect("load the native multipool cohort before its first claim"); + assert_eq!(before.replicas.len(), 6); + let now = unix_now(); + let expected = claim_scanner_pause_backlog_writer(&before.ledger, now).expect("the first writer generation"); + let acquire_timeout = ecstore_get_lock_acquire_timeout(); + assert_eq!(acquire_timeout, Duration::from_secs(5), "retain the production acquisition budget"); + let barrier = PutObjectCommitBarrier::install( + RUSTFS_META_BUCKET, + &SCANNER_PAUSE_BACKLOG_PATH, + PutObjectCommitPause::AfterNamespace, + ); + let writer_store = Arc::clone(&store); + let writer = tokio::spawn(async move { ScannerPauseBacklogController::claim(writer_store, now).await }); + barrier.wait_until_paused().await; + + // Every native multipool replica takes the same fixed outer lock. + // A parallel sibling would exhaust its unchanged budget while this + // write owns that lock. Keep real time and the native lock budget. + tokio::time::sleep(acquire_timeout + Duration::from_secs(1)).await; + assert!(!writer.is_finished(), "the first native write must remain behind the barrier"); + drop(barrier); + + let controller = writer + .await + .expect("the native writer task completes") + .expect("one publication must not time out acquiring locks held by its own replica writes"); + assert!(!controller.loaded.requires_reload); + assert_current_native_writer_ledger(&store, &expected).await; + drop(controller); + shutdown_native_retirement_store(store).await; + }); + } + #[test] #[serial_test::serial] fn native_writer_expansion_seed_failure_preserves_old_authority() { diff --git a/crates/scanner/src/storage_api.rs b/crates/scanner/src/storage_api.rs index c47ad1041..edeb0731f 100644 --- a/crates/scanner/src/storage_api.rs +++ b/crates/scanner/src/storage_api.rs @@ -150,7 +150,11 @@ pub(crate) mod owner { pub(crate) use super::{NativeScannerPauseBacklogWriteFault, SourceCleanupDeleteBarrier}; #[cfg(test)] - pub(crate) use rustfs_ecstore::api::set_disk::test_util::hold_namespace_commit as ecstore_hold_namespace_commit; + pub(crate) use rustfs_ecstore::api::set_disk::get_lock_acquire_timeout as ecstore_get_lock_acquire_timeout; + #[cfg(test)] + pub(crate) use rustfs_ecstore::api::set_disk::test_util::{ + PutObjectCommitBarrier, PutObjectCommitPause, hold_namespace_commit as ecstore_hold_namespace_commit, + }; pub(crate) use super::storage_contracts::{ HTTPPreconditions, HTTPRangeSpec, NS_SCANNER_PROTOCOL_VERSION, ObjectIO, ObjectOperations, ObjectToDelete,