fix(scanner): serialize pause backlog replica publication (#8264)

This commit is contained in:
Chris
2026-09-30 15:45:37 +08:00
committed by GitHub
parent d60dfbb826
commit fab2d7e144
2 changed files with 60 additions and 16 deletions
+55 -15
View File
@@ -1587,26 +1587,24 @@ where
.iter()
.filter_map(|replica| replica.revision.clone().map(|revision| (replica.id, revision)))
.collect::<HashMap<_, _>>();
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() {
+5 -1
View File
@@ -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,