diff --git a/crates/ecstore/src/ecstore_validation_blackbox.rs b/crates/ecstore/src/ecstore_validation_blackbox.rs index fbac4cd7f..15d30b81b 100644 --- a/crates/ecstore/src/ecstore_validation_blackbox.rs +++ b/crates/ecstore/src/ecstore_validation_blackbox.rs @@ -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"); diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 2b041ef13..bc6daedf6 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -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 { + 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>, instance_ctx: Arc, diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 305b9612a..956d89bd7 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -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; diff --git a/crates/heal/tests/mrf_pipeline_test.rs b/crates/heal/tests/mrf_pipeline_test.rs index 7cbf3d331..a93e13bb2 100644 --- a/crates/heal/tests/mrf_pipeline_test.rs +++ b/crates/heal/tests/mrf_pipeline_test.rs @@ -229,10 +229,42 @@ fn committed_manifest(owner: uuid::Uuid, sequence: u64, payload: &[u8]) -> 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; diff --git a/crates/scanner/src/scanner/backlog.rs b/crates/scanner/src/scanner/backlog.rs index b8f2eb314..3981df07e 100644 --- a/crates/scanner/src/scanner/backlog.rs +++ b/crates/scanner/src/scanner/backlog.rs @@ -1589,22 +1589,29 @@ where .collect::>(); // 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()