fix(storage): prevent metadata deadlocks and abandoned writes (#8268)

## Related Issues

N/A

## Summary of Changes

Reject system-bucket incarnation lookups before entering the pool metadata owner. Keep each selected scanner backlog publication cohort in one task so waiter cancellation cannot abandon its remaining serialized conditional writes. Bind repair and replay fixtures to persisted bucket incarnations.

## Verification

Head f6d44603fc received an approval from houseme in review 5365588011. This merge does not add a local runtime validation claim. Full main CI remains a separate publication gate.

## Impact

System metadata writes avoid recursive pool locking. Publication retains per-replica conditional writes and stops a cancelled caller from advancing to its next publication phase. Fixture changes supply the identities required by existing repair admission rules.

## Additional Notes

Reverting this change restores the previous behavior.
This commit is contained in:
Chris
2026-09-30 19:49:24 +08:00
committed by GitHub
parent eb350ad20d
commit 88d03d8199
5 changed files with 141 additions and 25 deletions
@@ -377,13 +377,19 @@ async fn blackbox_heal_requests_preserve_repair_scope() {
// Without a durable MRF consumer, partial PUTs fall back to the heal
// admission channel and wait for its receipt before acknowledging the write.
// Drive the test receiver alongside the PUT so neither waits on the other.
let (_put_dirs, put_set) = make_local_set_disks(4, 2).await;
let (_put_dirs, put_store) = crate::bucket::metadata_sys::test_support::isolated_store_over_temp_disks().await;
crate::bucket::metadata_sys::init_bucket_metadata_sys(put_store.clone(), Vec::new()).await;
let put_set = put_store.pools[0].disk_set[0].clone();
let put_bucket = "bb-put-partial-convergence";
let put_object = "object.bin";
put_set
put_store
.make_bucket(put_bucket, &MakeBucketOptions::default())
.await
.expect("PUT bucket should be created");
let put_incarnation = put_store
.bucket_incarnation_id_from_disk(put_bucket)
.await
.expect("PUT bucket should have a persisted incarnation");
let offline_disk = {
let mut disks = put_set.disks.write().await;
disks[0].take()
@@ -417,6 +423,7 @@ async fn blackbox_heal_requests_preserve_repair_scope() {
.to_string();
assert_eq!(request.object_version_id.as_deref(), Some(committed_version.as_str()));
assert_eq!(request.expected_bucket_incarnation_id, Some(put_incarnation));
assert_eq!(request.pool_index, Some(0));
assert_eq!(request.set_index, Some(0));
@@ -447,7 +454,7 @@ async fn blackbox_heal_requests_preserve_repair_scope() {
}
let healthy_bucket = "bb-put-healthy-convergence";
put_set
put_store
.make_bucket(healthy_bucket, &MakeBucketOptions::default())
.await
.expect("healthy PUT bucket should be created");
+33
View File
@@ -4809,6 +4809,10 @@ impl SetDisks {
/// Read the persisted bucket identity through this set's metadata owner.
/// Missing or non-authoritative legacy identities remain errors.
pub async fn bucket_incarnation_id_from_disk(&self, bucket: &str) -> Result<Uuid> {
if crate::bucket::utils::is_meta_bucketname(bucket) {
// Metadata writes can already hold the pool metadata write lock.
return Err(Error::other("system metadata bucket has no bucket incarnation"));
}
metadata_sys::get_bucket_incarnation_id_in(&self.ctx, bucket).await
}
@@ -7434,6 +7438,35 @@ mod tests {
make_test_set_disks_with_ctx(lockers, bootstrap_ctx()).await
}
#[tokio::test]
#[serial_test::serial]
async fn system_metadata_incarnation_lookup_does_not_reenter_pool_metadata() {
let (_temp_dirs, store, _other_store) =
crate::services::rebalance::test_three_pool_stores_with_isolated_node_contexts(None).await;
let _pool_meta_guard = store.pool_meta.write().await;
let set = &store.pools[0].disk_set[0];
for bucket in [RUSTFS_META_BUCKET, RUSTFS_META_TMP_BUCKET, crate::disk::MIGRATING_META_BUCKET] {
tokio::time::timeout(Duration::from_secs(30), set.bucket_incarnation_id_from_disk(bucket))
.await
.expect("system metadata identity lookup must not reacquire the held pool metadata lock")
.expect_err("system metadata buckets have no user bucket incarnation");
}
drop(_pool_meta_guard);
let bucket = "user-incarnation-boundary";
let incarnation = Uuid::new_v4();
crate::bucket::metadata::save_bucket_incarnation(Arc::clone(&store), bucket, incarnation)
.await
.expect("persist the user bucket identity through the metadata owner");
assert_eq!(
set.bucket_incarnation_id_from_disk(bucket)
.await
.expect("user bucket identities must still load from the metadata owner"),
incarnation,
);
}
async fn make_test_set_disks_with_ctx(
lockers: Vec<Arc<dyn LockClient>>,
instance_ctx: Arc<InstanceContext>,
+28 -2
View File
@@ -3194,9 +3194,35 @@ mod tests {
replay_intent.kind = MrfKind::PartialWrite;
replay_intent.version_id = None;
let replay_payload = encoded_payload(&replay_intent);
snapshot::publish_committed_snapshot(&disks, replay_owner, 11, &replay_payload, config.journal_max_bytes)
let source_incarnation = storage
.mrf_bucket_incarnation_id(bucket)
.await
.expect("publish committed replay checkpoint");
.expect("read persisted replay source identity")
.expect("replay source bucket has a persisted incarnation");
let lifecycle_limit = config.journal_max_bytes.saturating_mul(4);
let lifecycle = encode_mrf_lifecycle_checkpoint(
replay_owner,
11,
vec![ResponsibilityCheckpoint {
intent_digest: intent_digest(&replay_intent).expect("fixture intent has a canonical digest"),
responsibility_id: Uuid::new_v4(),
source_bucket_incarnation_id: Some(source_incarnation),
last_operator_acceptance: None,
state: partial_write::ResponsibilityState::Active,
}],
lifecycle_limit,
)
.expect("encode generation-bound replay responsibility");
snapshot::publish_committed_snapshot_with_companion(
&disks,
replay_owner,
11,
&replay_payload,
config.journal_max_bytes,
Some((&MRF_LIFECYCLE_PATHS, &lifecycle, lifecycle_limit)),
)
.await
.expect("publish committed replay checkpoint with its source identity");
let mut queue = MrfQueue::new(config.queue_capacity, config.journal_max_bytes);
let mut backoff_until = None;
+47 -4
View File
@@ -229,10 +229,42 @@ fn committed_manifest(owner: uuid::Uuid, sequence: u64, payload: &[u8]) -> Vec<u
manifest
}
fn write_committed_snapshot_to_disks(disk_paths: &[std::path::PathBuf], sequence: u64, payload: &[u8]) {
let manifest = committed_manifest(uuid::Uuid::new_v4(), sequence, payload);
fn write_generation_bound_snapshot_to_disks(
disk_paths: &[PathBuf],
sequence: u64,
payload: &[u8],
source_incarnation: uuid::Uuid,
partial_records: &[&[u8]],
) {
assert!(!source_incarnation.is_nil(), "replay source must come from the persisted bucket");
let owner = uuid::Uuid::new_v4();
let records = partial_records
.iter()
.map(|record| {
serde_json::json!({
"intent_digest": Sha256::digest(record).to_vec(),
"responsibility_id": uuid::Uuid::new_v4(),
"source_bucket_incarnation_id": source_incarnation,
"last_operator_acceptance": null,
"state": { "state": "active" },
})
})
.collect::<Vec<_>>();
let lifecycle = serde_json::to_vec(&serde_json::json!({
"format_version": 1,
"checkpoint_owner": owner,
"checkpoint_sequence": sequence,
"records": records,
}))
.expect("encode source-bound lifecycle fixture");
let mut companion = b"RFMFLC01".to_vec();
companion.push(1);
companion.extend_from_slice(&u64::try_from(lifecycle.len()).expect("fixture length fits").to_le_bytes());
companion.extend_from_slice(&Sha256::digest(&lifecycle));
companion.extend_from_slice(&lifecycle);
write_journal_path_to_disks(disk_paths, ".heal-mrf-lifecycle.0.bin", &companion);
write_journal_path_to_disks(disk_paths, COMMITTED_PAYLOAD_REL, payload);
write_journal_path_to_disks(disk_paths, COMMITTED_MANIFEST_REL, &manifest);
write_journal_path_to_disks(disk_paths, COMMITTED_MANIFEST_REL, &committed_manifest(owner, sequence, payload));
}
fn journal_exists_on_all_disks(disk_paths: &[std::path::PathBuf], relative_path: &str) -> bool {
@@ -423,7 +455,12 @@ async fn committed_snapshot_replay_takes_precedence_over_stale_legacy_mirror() {
let committed = scoped_journal_record(3, "committed-bucket", "committed-object", Some([9u8; 16]), 0, 0, 0);
let stale_legacy = journal_record(1, "legacy-bucket", "legacy-object", None, 0);
write_committed_snapshot_to_disks(&disk_paths, 7, &committed);
let incarnation = storage
.mrf_bucket_incarnation_id("committed-bucket")
.await
.expect("read committed source incarnation")
.expect("committed source bucket has a persisted incarnation");
write_generation_bound_snapshot_to_disks(&disk_paths, 7, &committed, incarnation, &[&committed]);
write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &stale_legacy);
write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &stale_legacy);
@@ -562,6 +599,12 @@ async fn authoritative_journal_replay_preserves_kind_and_scope_identity() {
let stale_legacy = journal_record(3, "identity-bucket", "stale-legacy-object", None, 0);
write_journal_path_to_disks(&disk_paths, SCOPED_JOURNAL_REL, &authoritative);
write_journal_path_to_disks(&disk_paths, JOURNAL_REL, &stale_legacy);
let incarnation = storage
.mrf_bucket_incarnation_id("identity-bucket")
.await
.expect("read authoritative source incarnation")
.expect("authoritative source bucket has a persisted incarnation");
write_generation_bound_snapshot_to_disks(&disk_paths, 1, &authoritative, incarnation, &[&first_partial, &second_partial]);
let manager = make_manager(storage);
let replayed = mrf_queue::replay_journal_once(&manager).await;
+23 -16
View File
@@ -1589,22 +1589,29 @@ where
.collect::<HashMap<_, _>>();
// 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 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()),
None => Err("replica revision is unavailable".to_string()),
};
results.push((id, result));
}
// Own the whole selected cohort so canceling the waiter cannot abandon
// replicas that have not started their serialized write yet.
let results = tokio::spawn(async move {
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 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()),
None => Err("replica revision is unavailable".to_string()),
};
results.push((id, result));
}
results
})
.await
.map_err(|err| format!("scanner pause backlog publication owner failed: {err}"))?;
let failures = results
.iter()