From cc5487e7de309e58b7108f839f8b0ac7f9b80072 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 16:57:13 +0800 Subject: [PATCH 01/12] heal: verify admin recreate pool metadata (#7474) Co-authored-by: zhi22915 --- crates/heal/src/heal/erasure_healer.rs | 91 +++++++++++++++++-- crates/heal/src/heal/task.rs | 2 +- crates/heal/src/heal/task/heal_bucket.rs | 79 ++++++++++++++++ crates/heal/src/heal/task/heal_erasure_set.rs | 5 + crates/heal/src/heal/task/heal_object.rs | 5 - crates/heal/src/heal/task/tests.rs | 88 +++++++++++++++++- 6 files changed, 252 insertions(+), 18 deletions(-) diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index dbd62e422..9cb5071bc 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -113,6 +113,7 @@ pub struct ErasureSetHealer { heal_opts: HealOpts, source: HealRequestSource, target_endpoints: Arc<[String]>, + pool_metadata_target_endpoints: Arc<[String]>, replacement_task_id: Option, replacement_target_identities: Option>, mainline_pacer: Option>, @@ -362,6 +363,7 @@ impl ErasureSetHealer { heal_opts, source, target_endpoints: Vec::new().into(), + pool_metadata_target_endpoints: Vec::new().into(), replacement_task_id: None, replacement_target_identities: None, mainline_pacer: None, @@ -385,6 +387,13 @@ impl ErasureSetHealer { self } + pub(crate) fn with_pool_metadata_targets(mut self, mut target_endpoints: Vec) -> Self { + target_endpoints.sort_unstable(); + target_endpoints.dedup(); + self.pool_metadata_target_endpoints = target_endpoints.into(); + self + } + pub(crate) fn with_replacement_identity_fence( mut self, replacement_target_identities: Option>, @@ -948,10 +957,16 @@ impl ErasureSetHealer { resume_manager: &ResumeManager, checkpoint_manager: &CheckpointManager, ) -> Result<()> { - if self.replacement_task_id.is_none() { + let target_endpoints = if self.pool_metadata_target_endpoints.is_empty() { + self.target_endpoints.as_ref() + } else { + self.pool_metadata_target_endpoints.as_ref() + }; + let target_scoped_recreate = !self.heal_opts.dry_run && self.heal_opts.recreate && !target_endpoints.is_empty(); + if self.replacement_task_id.is_none() && !target_scoped_recreate { return Ok(()); } - if self.target_endpoints.is_empty() { + if target_endpoints.is_empty() { return Err(Error::TaskExecutionFailed { message: "Replacement pool metadata heal requires target endpoints".to_string(), }); @@ -978,17 +993,11 @@ impl ErasureSetHealer { .heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, &self.heal_opts) .await { - Ok((result, None)) if target_outcomes_complete(&result, &self.target_endpoints) => { + Ok((result, None)) if target_outcomes_complete(&result, target_endpoints) => { let object_size = result_object_size_u64(&result); match self .storage - .replacement_targets_have_version( - RUSTFS_META_BUCKET, - POOL_META_NAME, - None, - &self.heal_opts, - &self.target_endpoints, - ) + .replacement_targets_have_version(RUSTFS_META_BUCKET, POOL_META_NAME, None, &self.heal_opts, target_endpoints) .await { Ok(true) => (object_size, Ok(())), @@ -2762,6 +2771,68 @@ mod resume_loop_tests { assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); } + #[tokio::test] + async fn admin_recreate_target_heals_pool_metadata_before_completion() { + let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; + let healer = ErasureSetHealer::new( + env.storage.clone(), + Arc::new(RwLock::new(HealProgress::new())), + CancellationToken::new(), + env.healer.disk.clone(), + HealOpts { + recreate: true, + pool: Some(0), + set: Some(0), + ..Default::default() + }, + HealRequestSource::Admin, + ) + .with_pool_metadata_targets(vec!["replacement-a".to_string()]); + env.storage + .set_result(POOL_META_NAME, None, replacement_target_ok_result("replacement-a", POOL_META_NAME)); + + healer + .execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect("admin recreate should heal and verify pool metadata"); + + assert!(env.resume.get_state().await.completed); + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + } + + #[tokio::test] + async fn admin_recreate_pool_metadata_readback_failure_keeps_resume_state() { + let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; + let healer = ErasureSetHealer::new( + env.storage.clone(), + Arc::new(RwLock::new(HealProgress::new())), + CancellationToken::new(), + env.healer.disk.clone(), + HealOpts { + recreate: true, + pool: Some(0), + set: Some(0), + ..Default::default() + }, + HealRequestSource::Admin, + ) + .with_pool_metadata_targets(vec!["replacement-a".to_string()]); + env.storage + .set_result(POOL_META_NAME, None, replacement_target_ok_result("replacement-a", POOL_META_NAME)); + env.storage.set_replacement_commit_evidence(POOL_META_NAME, None, false); + + let error = healer + .execute_heal_with_resume(&["b".to_string()], "pool_0_set_0", &env.resume, &env.checkpoint) + .await + .expect_err("unconfirmed admin recreate pool metadata must keep the set incomplete"); + + assert!(error.to_string().contains("Erasure set heal incomplete")); + let state = env.resume.get_state().await; + assert!(!state.completed); + assert_eq!(state.retry_count, 1); + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + } + #[tokio::test] async fn retry_exhaustion_keeps_resume_artifacts_for_recovery() { let env = make_env().await; diff --git a/crates/heal/src/heal/task.rs b/crates/heal/src/heal/task.rs index 2c171e76e..816f22afd 100644 --- a/crates/heal/src/heal/task.rs +++ b/crates/heal/src/heal/task.rs @@ -45,7 +45,7 @@ use tokio::sync::RwLock; use tracing::{debug, error, info, warn}; use uuid::Uuid; -use super::{BUCKET_META_PREFIX, DATA_USAGE_CACHE_NAME, RUSTFS_META_BUCKET}; +use super::{BUCKET_META_PREFIX, DATA_USAGE_CACHE_NAME, POOL_META_NAME, RUSTFS_META_BUCKET}; #[cfg(test)] pub(crate) struct OutcomeFinishTestHook { diff --git a/crates/heal/src/heal/task/heal_bucket.rs b/crates/heal/src/heal/task/heal_bucket.rs index d7465fdb5..cec17886a 100644 --- a/crates/heal/src/heal/task/heal_bucket.rs +++ b/crates/heal/src/heal/task/heal_bucket.rs @@ -340,9 +340,88 @@ impl HealTask { return Err(self.record_batch_failure(failure).await); } + if self.options.recreate_missing && !self.options.dry_run { + self.heal_cluster_pool_metadata().await?; + } + Ok(()) } + async fn heal_cluster_pool_metadata(&self) -> Result<()> { + let heal_opts = HealOpts { + recursive: false, + dry_run: self.options.dry_run, + remove: false, + recreate: self.options.recreate_missing, + scan_mode: self.options.scan_mode, + update_parity: self.options.update_parity, + no_lock: self.options.no_lock, + read_repair: false, + pool: self.options.pool_index, + set: self.options.set_index, + }; + + let heal_result = self + .await_with_control(self.storage.heal_object(RUSTFS_META_BUCKET, POOL_META_NAME, None, &heal_opts)) + .await; + match heal_result { + Ok((result, None)) => { + debug!( + target: "rustfs::heal::task", + event = EVENT_HEAL_BUCKET_RESULT, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_TASK, + task_id = %self.id, + bucket = RUSTFS_META_BUCKET, + object = POOL_META_NAME, + drives_healed = result.drives_healed(), + drives_total = result.drives_reported(), + result = "pool_metadata_ok", + "Heal cluster pool metadata repaired" + ); + self.record_result_item(result).await; + Ok(()) + } + Ok((result, Some(err))) => { + self.record_result_item(result).await; + warn!( + target: "rustfs::heal::task", + event = EVENT_HEAL_BUCKET_RESULT, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_TASK, + task_id = %self.id, + bucket = RUSTFS_META_BUCKET, + object = POOL_META_NAME, + result = "pool_metadata_failed", + error = %err, + "Heal cluster pool metadata failed" + ); + Err(Error::TaskExecutionFailed { + message: format!("Failed to heal cluster pool metadata: {err}"), + }) + } + Err(Error::TaskCancelled) => Err(Error::TaskCancelled), + Err(Error::TaskTimeout) => Err(Error::TaskTimeout), + Err(err) => { + warn!( + target: "rustfs::heal::task", + event = EVENT_HEAL_BUCKET_RESULT, + component = LOG_COMPONENT_HEAL, + subsystem = LOG_SUBSYSTEM_TASK, + task_id = %self.id, + bucket = RUSTFS_META_BUCKET, + object = POOL_META_NAME, + result = "pool_metadata_failed", + error = %err, + "Heal cluster pool metadata failed" + ); + Err(Error::TaskExecutionFailed { + message: format!("Failed to heal cluster pool metadata: {err}"), + }) + } + } + } + pub(super) async fn heal_prefix(&self, bucket: &str, prefix: &str) -> Result<()> { debug!( target: "rustfs::heal::task", diff --git a/crates/heal/src/heal/task/heal_erasure_set.rs b/crates/heal/src/heal/task/heal_erasure_set.rs index 76b65a41a..2a6ccc494 100644 --- a/crates/heal/src/heal/task/heal_erasure_set.rs +++ b/crates/heal/src/heal/task/heal_erasure_set.rs @@ -451,6 +451,11 @@ impl HealTask { self.source, ) .with_replacement_targets(replacement_targets, is_auto_replacement.then(|| self.id.clone())) + .with_pool_metadata_targets(if self.options.recreate_missing && !self.options.dry_run { + self.heal_endpoints.clone() + } else { + Vec::new() + }) .with_replacement_identity_fence(replacement_target_identities.clone()) .with_mainline_pacer(self.mainline_pacer.clone()); diff --git a/crates/heal/src/heal/task/heal_object.rs b/crates/heal/src/heal/task/heal_object.rs index 2265d1919..c7635a79f 100644 --- a/crates/heal/src/heal/task/heal_object.rs +++ b/crates/heal/src/heal/task/heal_object.rs @@ -162,11 +162,6 @@ impl HealTask { pool: self.options.pool_index, set: self.options.set_index, }; - let expected_bucket_incarnation_id = self.storage.bucket_incarnation_id(bucket).await?; - let mut expected_identity = - self.outcome_identity(bucket, object, version_id, self.options.pool_index, self.options.set_index); - expected_identity.bucket_incarnation_id = expected_bucket_incarnation_id; - let mut expected_identity = self.outcome_identity(bucket, object, version_id, self.options.pool_index, self.options.set_index); expected_identity.bucket_incarnation_id = self.outcome_bucket_incarnation_id(bucket, self.options.dry_run).await?; diff --git a/crates/heal/src/heal/task/tests.rs b/crates/heal/src/heal/task/tests.rs index 9981d73d3..63376d08e 100644 --- a/crates/heal/src/heal/task/tests.rs +++ b/crates/heal/src/heal/task/tests.rs @@ -66,7 +66,7 @@ mod canonical_outcome { assert_eq!(task.get_progress().await.objects_scanned, 2); assert_eq!( storage.heal_object_calls.lock().expect("object calls").as_slice(), - ["object-a", "object-b"] + ["object-a", "object-b", POOL_META_NAME] ); assert_eq!( storage.listing_tokens.lock().expect("listing tokens").as_slice(), @@ -2582,11 +2582,95 @@ async fn test_cluster_heal_visits_bucket_objects() { assert_eq!( storage.healed_objects.lock().unwrap().as_slice(), - ["object-a".to_string(), "object-b".to_string()] + ["object-a".to_string(), "object-b".to_string(), POOL_META_NAME.to_string()] ); assert!(matches!(task.get_status().await, HealTaskStatus::Completed)); } +#[tokio::test] +async fn cluster_recreate_heals_pool_metadata_after_user_buckets() { + let storage = Arc::new(MockStorage::default()); + let request = HealRequest::new( + HealType::Cluster, + HealOptions { + recursive: true, + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("cluster recreate heal should include pool metadata"); + + assert_eq!( + storage.heal_object_calls.lock().expect("object calls").as_slice(), + ["object-a".to_string(), "object-b".to_string(), POOL_META_NAME.to_string()] + ); + let opts = storage.object_heal_opts.lock().expect("object opts"); + assert!(opts.last().expect("pool metadata opts").recreate); +} + +#[tokio::test] +async fn cluster_recreate_fails_when_pool_metadata_heal_fails() { + let storage = Arc::new(MockStorage::default()); + storage.heal_object_outcomes.lock().expect("object outcomes").insert( + POOL_META_NAME.to_string(), + VecDeque::from([MockHealObjectOutcome::ErrOther("pool metadata missing")]), + ); + let request = HealRequest::new( + HealType::Cluster, + HealOptions { + recursive: true, + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + let task = HealTask::from_request(request, storage.clone()); + + let err = task + .execute() + .await + .expect_err("cluster recreate heal must not hide pool metadata failure"); + + assert!(matches!(err, Error::TaskExecutionFailed { .. })); + assert_eq!( + storage.heal_object_calls.lock().expect("object calls").as_slice(), + ["object-a".to_string(), "object-b".to_string(), POOL_META_NAME.to_string()] + ); +} + +#[tokio::test] +async fn cluster_dry_run_does_not_heal_pool_metadata() { + let storage = Arc::new(MockStorage::default()); + let request = HealRequest::new( + HealType::Cluster, + HealOptions { + recursive: true, + dry_run: true, + recreate_missing: true, + timeout: None, + ..Default::default() + }, + HealPriority::Normal, + ); + let task = HealTask::from_request(request, storage.clone()); + + task.execute() + .await + .expect("dry-run cluster heal should preserve existing coverage"); + + assert_eq!( + storage.heal_object_calls.lock().expect("object calls").as_slice(), + ["object-a".to_string(), "object-b".to_string()] + ); +} + #[tokio::test] async fn object_heal_skips_dangling_delete_grace_without_failing_task() { let storage = Arc::new(MockStorage { From 5dc0e3b40258876a8cb995249aa84a3c1dc91761 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 16:59:17 +0800 Subject: [PATCH 02/12] heal: preserve merged result for duplicate submits Keep same-request-id replay receipts accepted for the receipt API, but preserve the legacy submit_heal_request duplicate admission result as Merged. Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- crates/heal/src/heal/manager.rs | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/crates/heal/src/heal/manager.rs b/crates/heal/src/heal/manager.rs index 6175d9583..732778d15 100644 --- a/crates/heal/src/heal/manager.rs +++ b/crates/heal/src/heal/manager.rs @@ -1527,7 +1527,7 @@ impl HealManager { request: HealRequest, preserve_alias: bool, ) -> Result { - self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, preserve_alias, None) + self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, preserve_alias, true, None) .await } @@ -1563,7 +1563,7 @@ impl HealManager { request: HealRequest, mrf_notice_target: MrfRepairNoticeTarget, ) -> Result { - self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, true, Some(mrf_notice_target)) + self.submit_heal_request_with_receipt_alias_and_mrf_notice(request, true, true, Some(mrf_notice_target)) .await } @@ -1583,6 +1583,7 @@ impl HealManager { &self, request: HealRequest, preserve_alias: bool, + accept_same_request_id_replay: bool, mrf_notice_target: Option, ) -> Result { let admission_start = Instant::now(); @@ -1661,7 +1662,11 @@ impl HealManager { }); if let Some((matches_existing, duplicate_state)) = request_id_admission { let admission = if matches_existing { - HealAdmissionResult::Accepted + if accept_same_request_id_replay { + HealAdmissionResult::Accepted + } else { + Self::duplicate_admission_for_request(&request, &config) + } } else { HealAdmissionResult::Dropped(HealAdmissionDropReason::AlreadyRunning) }; @@ -1908,7 +1913,10 @@ impl HealManager { /// Submit heal request. pub async fn submit_heal_request(&self, request: HealRequest) -> Result { - Ok(self.submit_heal_request_with_receipt_and_alias(request, true).await?.result) + Ok(self + .submit_heal_request_with_receipt_alias_and_mrf_notice(request, true, false, None) + .await? + .result) } /// Get task status From c0754f5b1c677ff128574ee0c95d22b56cd43eda Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 17:26:11 +0800 Subject: [PATCH 03/12] test(heal): preserve EC84 restart semantics (#7478) Keep the EC8+4 background restart lane on the graceful-stop path and assert clean-restart marker absence only for restart scenarios. This prevents the hard evidence gate from silently exercising the crash path when it claims restart coverage. Co-authored-by: zhi22915 --- crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs index 0b6e30bd6..f8763d482 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -1567,7 +1567,10 @@ mod tests { "Restored target endpoint forwarding" ); } else { - if scenario == InterruptionScenario::BackgroundTargetRestart { + if matches!( + scenario, + InterruptionScenario::BackgroundTargetRestart | InterruptionScenario::BackgroundTargetRestartEc84 + ) { cluster.stop_node_gracefully(interruption_node).await?; } else { cluster.stop_node(interruption_node)?; @@ -1589,7 +1592,10 @@ mod tests { if background_enabled { let marker_exists = unclean_shutdown_marker.is_file(); unclean_shutdown_marker_observed = Some(marker_exists); - let expected_marker = !matches!(scenario, InterruptionScenario::BackgroundTargetRestart); + let expected_marker = matches!( + scenario, + InterruptionScenario::BackgroundTargetCrash | InterruptionScenario::BackgroundTargetCrashEc84 + ); assert!( marker_exists == expected_marker, "background restart/crash lane observed unexpected unclean-shutdown marker state" From 7b3dad6bae0b8492d05f6f3b3e6805ea3f28a0d1 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 17:27:47 +0800 Subject: [PATCH 04/12] test(heal): cover MRF snapshot torn successor recovery Co-Authored-By: heihutu Co-Authored-By: zhi22915 --- crates/heal/src/heal/mrf_queue/snapshot.rs | 77 ++++++++++++++++++++++ 1 file changed, 77 insertions(+) diff --git a/crates/heal/src/heal/mrf_queue/snapshot.rs b/crates/heal/src/heal/mrf_queue/snapshot.rs index 43fc7372b..0921c043f 100644 --- a/crates/heal/src/heal/mrf_queue/snapshot.rs +++ b/crates/heal/src/heal/mrf_queue/snapshot.rs @@ -1200,6 +1200,83 @@ mod tests { assert_eq!(recovered.manifest.sequence, 1); } + #[tokio::test] + async fn manifest_cas_failure_after_payload_write_keeps_previous_anchor() { + let root = TempDir::new().expect("test directory"); + let store = disk(&root, "disk").await; + let owner = Uuid::new_v4(); + let old = payload("old"); + let next = payload("next"); + let damaged_manifest = b"damaged successor manifest".to_vec(); + commit(&store, 0, owner, 1, &old).await; + + let expected_manifest = EcstoreDiskAPI::read_all(store.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[1]) + .await + .ok(); + assert_eq!( + cas_replace(&store, PAYLOAD_PATHS[1], &next, 4096) + .await + .expect("successor payload CAS"), + EcstoreConditionalFileUpdate::Updated + ); + install(&store, MANIFEST_PATHS[1], &damaged_manifest).await; + + let manifest_update = cas_replace_expected(&store, MANIFEST_PATHS[1], expected_manifest, &manifest(owner, 2, &next)) + .await + .expect("successor manifest CAS"); + assert_eq!(manifest_update, EcstoreConditionalFileUpdate::Mismatch); + + let reopened = disk(&root, "disk").await; + let recovered = read_committed(std::slice::from_ref(&reopened), 4096) + .await + .expect("read committed snapshot after failed successor CAS") + .expect("previous committed anchor"); + assert_eq!(recovered.sequence(), 1); + assert_eq!(recovered.slot(), 0); + assert_eq!(recovered.payload(), old.as_slice()); + assert_eq!( + EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, PAYLOAD_PATHS[1]) + .await + .expect("successor payload remains non-authoritative") + .as_ref(), + next.as_slice() + ); + assert_eq!( + EcstoreDiskAPI::read_all(reopened.as_ref(), RUSTFS_META_BUCKET, MANIFEST_PATHS[1]) + .await + .expect("failed successor manifest retained") + .as_ref(), + damaged_manifest.as_slice() + ); + } + + #[tokio::test] + async fn torn_successor_on_one_replica_does_not_hide_previous_anchor_on_peer() { + let root = TempDir::new().expect("test directory"); + let first = disk(&root, "first").await; + let second = disk(&root, "second").await; + let owner = Uuid::new_v4(); + let old = payload("old"); + let next = payload("next"); + let damaged_manifest = b"damaged successor manifest".to_vec(); + commit(&first, 0, owner, 1, &old).await; + commit(&second, 0, owner, 1, &old).await; + install(&first, PAYLOAD_PATHS[1], &next).await; + install(&first, MANIFEST_PATHS[1], &damaged_manifest).await; + + let mut stats = SnapshotReadStats::default(); + let recovered = read_committed_with_stats(&[first, second], 4096, Some(&mut stats)) + .await + .expect("read committed snapshot across torn successor") + .expect("previous committed anchor"); + + assert_eq!(recovered.sequence(), 1); + assert_eq!(recovered.payload(), old.as_slice()); + assert_eq!(stats.file_reads, 5); + assert_eq!(stats.bytes_read, (MANIFEST_LEN * 2) + (old.len() * 2) + damaged_manifest.len()); + assert_eq!(stats.peak_file_bytes, old.len().max(next.len()).max(MANIFEST_LEN)); + } + #[tokio::test] async fn manifest_cas_publication_transitions_from_legacy_without_losing_anchor() { let root = TempDir::new().expect("test directory"); From 2ba7f95547aad01c8947d37d134e7e6218c2f209 Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 18:15:46 +0800 Subject: [PATCH 05/12] test(heal): write EC84 distributed restart oracle (#7481) Bind the distributed EC8+4 restart evidence lane to a scanner/heal oracle artifact so release validation can consume the real nextest run instead of accepting only a passing test. Require the registry to assert 8+4 erasure geometry for the three-node, four-drive case. Co-authored-by: zhi22915 --- .config/scanner-heal-required-tests.json | 2 + crates/e2e_test/src/distributed/heal_test.rs | 183 ++++++++++++++++++- 2 files changed, 184 insertions(+), 1 deletion(-) diff --git a/.config/scanner-heal-required-tests.json b/.config/scanner-heal-required-tests.json index ac6583d0c..2f265da53 100644 --- a/.config/scanner-heal-required-tests.json +++ b/.config/scanner-heal-required-tests.json @@ -43,6 +43,8 @@ "min_objects": 5, "max_objects": 5, "topology": {"nodes": 3, "drives_per_node": 4}, + "erasure": {"data_blocks": 8, "parity_blocks": 4}, + "erasure_set_drive_count": 12, "scope": "3-node x 4-drive single-set EC8+4, graceful target restart, preformatted replacement drive, exact unversioned S3 bodies and physical target shards; not mixed-version, multi-pool or long-window ABBA." }, "background-target-restart-ec8-4": { diff --git a/crates/e2e_test/src/distributed/heal_test.rs b/crates/e2e_test/src/distributed/heal_test.rs index 01d3f099a..f0e0516cd 100644 --- a/crates/e2e_test/src/distributed/heal_test.rs +++ b/crates/e2e_test/src/distributed/heal_test.rs @@ -12,12 +12,18 @@ // See the License for the specific language governing permissions and // limitations under the License. -use super::harness::{DistCluster, DistLayout, TestResult, assert_inventory, payload_for, put_object, unique_bucket, wait_until}; +use super::harness::{ + DistCluster, DistLayout, TestResult, assert_inventory, get_object_bytes, payload_for, put_object, sha256_hex, unique_bucket, + wait_until, +}; use crate::chaos::{VersionShardCensus, census_object_version_on_disk, signed_admin_post}; use crate::common::init_logging; use aws_sdk_s3::Client; use aws_sdk_s3::primitives::ByteStream; +use serde_json::Value; +use sha2::{Digest, Sha256}; use std::collections::{BTreeMap, HashSet}; +use std::io::{Read, Write}; use std::path::{Path, PathBuf}; use std::time::Duration; @@ -25,6 +31,8 @@ const EC84_NODE_COUNT: usize = 3; const EC84_DRIVES_PER_NODE: usize = 4; const EC84_DATA_BLOCKS: usize = 8; const EC84_PARITY_BLOCKS: usize = 4; +const EC84_TARGET_DRIVE_RESTART_CASE: &str = "ec84-target-drive-restart"; +const EC84_TARGET_DRIVE_RESTART_ORACLE: &str = "ec84-target-drive-restart.json"; #[derive(Clone)] struct ExpectedShard { @@ -33,6 +41,79 @@ struct ExpectedShard { baseline: VersionShardCensus, } +struct ScannerHealEvidenceContext { + directory: PathBuf, + run: Value, +} + +fn file_sha256(path: &Path) -> TestResult { + let mut file = std::fs::File::open(path)?; + let mut digest = Sha256::new(); + let mut buffer = [0_u8; 64 * 1024]; + loop { + let read = file.read(&mut buffer)?; + if read == 0 { + break; + } + digest.update(&buffer[..read]); + } + Ok(digest.finalize().iter().map(|byte| format!("{byte:02x}")).collect()) +} + +fn compiled_test_identity() -> Value { + serde_json::json!({ + "source_revision": env!("RUSTFS_E2E_BUILD_COMMIT"), + "dirty": env!("RUSTFS_E2E_BUILD_DIRTY") != "false", + "lock_blob": env!("RUSTFS_E2E_BUILD_LOCK"), + "features": env!("RUSTFS_E2E_BUILD_FEATURES"), + "target": env!("RUSTFS_E2E_BUILD_TARGET"), + "profile": env!("RUSTFS_E2E_BUILD_PROFILE"), + "rustflags_hex": env!("RUSTFS_E2E_BUILD_RUSTFLAGS_HEX"), + }) +} + +fn string_field<'a>(value: &'a Value, path: &str) -> TestResult<&'a str> { + let mut current = value; + for segment in path.split('.') { + current = current + .get(segment) + .ok_or_else(|| format!("scanner/heal run receipt missing {path}"))?; + } + current + .as_str() + .filter(|text| !text.is_empty()) + .ok_or_else(|| format!("scanner/heal run receipt has invalid {path}").into()) +} + +fn scanner_heal_evidence_context() -> TestResult> { + let Some(directory) = std::env::var_os("RUSTFS_SCANNER_HEAL_RUN_DIR") else { + return Ok(None); + }; + let directory = PathBuf::from(directory); + let receipt = directory.join("run.json"); + if receipt.metadata()?.len() > 1024 * 1024 { + return Err("oversized scanner/heal execution receipt".into()); + } + let run: Value = serde_json::from_slice(&std::fs::read(receipt)?)?; + let built = compiled_test_identity(); + for key in ["source_revision", "dirty", "lock_blob", "features"] { + if built[key] != run["test_build"][key] { + return Err(format!("compiled test identity differs for {key}").into()); + } + } + let binary_path = PathBuf::from(string_field(&run, "binary.path")?); + if file_sha256(&binary_path)? != string_field(&run, "binary.sha256")? { + return Err("server binary must match the run receipt".into()); + } + if file_sha256(&std::env::current_exe()?)? != string_field(&run, "test_binary.sha256")? { + return Err("test executable must match the run receipt".into()); + } + if directory.join(EC84_TARGET_DRIVE_RESTART_ORACLE).exists() { + return Err("scanner/heal oracle already exists; create a new execution receipt".into()); + } + Ok(Some(ScannerHealEvidenceContext { directory, run })) +} + fn assert_ec84_geometry(census: &VersionShardCensus, key: &str) -> TestResult { if census.data_blocks != Some(EC84_DATA_BLOCKS) || census.parity_blocks != Some(EC84_PARITY_BLOCKS) { return Err(format!("object {key} did not use EC8+4 geometry: {census:?}").into()); @@ -49,6 +130,76 @@ fn assert_ec84_geometry(census: &VersionShardCensus, key: &str) -> TestResult { Ok(()) } +async fn write_scanner_heal_evidence( + context: ScannerHealEvidenceContext, + dist: &DistCluster, + bucket: &str, + expected: &[ExpectedShard], + outage_key: &str, + outage_body: &[u8], + replaced_drive: &Path, + pid_before: u32, + pid_after: u32, + node_listings: Vec>, +) -> TestResult { + let verifier = dist.client(0)?; + let mut objects = Vec::new(); + for item in expected { + let actual = get_object_bytes(&verifier, bucket, &item.key).await?; + let physical = census_object_version_on_disk(replaced_drive, bucket, &item.key, None)?; + objects.push(serde_json::json!({ + "key": item.key, + "version_id": null, + "expected_bytes": item.body.len(), + "actual_bytes": actual.len(), + "expected_sha256": sha256_hex(&item.body), + "actual_sha256": sha256_hex(&actual), + "expected_physical": item.baseline, + "physical": physical, + })); + } + let actual = get_object_bytes(&verifier, bucket, outage_key).await?; + let physical = census_object_version_on_disk(replaced_drive, bucket, outage_key, None)?; + objects.push(serde_json::json!({ + "key": outage_key, + "version_id": null, + "expected_bytes": outage_body.len(), + "actual_bytes": actual.len(), + "expected_sha256": sha256_hex(outage_body), + "actual_sha256": sha256_hex(&actual), + "expected_physical": null, + "physical": physical, + })); + + let evidence = serde_json::json!({ + "schema": 1, + "case": EC84_TARGET_DRIVE_RESTART_CASE, + "evidence": "process-restart", + "run_id": string_field(&context.run, "run_id")?, + "source_revision": string_field(&context.run, "source_revision")?, + "test_build": compiled_test_identity(), + "binary_sha256": string_field(&context.run, "binary.sha256")?, + "test_binary_sha256": string_field(&context.run, "test_binary.sha256")?, + "topology": {"nodes": EC84_NODE_COUNT, "drives_per_node": EC84_DRIVES_PER_NODE}, + "pid_before": pid_before, + "pid_after": pid_after, + "unclean_shutdown_marker": false, + "objects": objects, + "node_listings": node_listings, + }); + let data = serde_json::to_vec(&evidence)?; + if data.len() > 1024 * 1024 { + return Err("scanner/heal oracle exceeds the 1 MiB artifact budget".into()); + } + let mut output = std::fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(context.directory.join(EC84_TARGET_DRIVE_RESTART_ORACLE))?; + output.write_all(&data)?; + output.sync_all()?; + Ok(()) +} + fn assert_replaced_drive_empty(drive: &Path, bucket: &str, keys: &[String]) -> TestResult { for key in keys { let census = census_object_version_on_disk(drive, bucket, key, None)?; @@ -87,6 +238,7 @@ async fn put_large_inventory(client: &Client, bucket: &str) -> TestResult TestResult { init_logging(); + let evidence_context = scanner_heal_evidence_context()?; let mut dist = DistCluster::start_with_env( DistLayout::ThreeByFourEc84, &[ @@ -116,6 +268,11 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res let format_path = replaced_drive.join(".rustfs.sys").join("format.json"); let format_json = std::fs::read(&format_path)?; + let target_pid_before = dist.cluster.nodes[replaced_node] + .process + .as_ref() + .ok_or("target process is absent before graceful restart")? + .id(); dist.cluster.stop_node_gracefully(replaced_node).await?; let retired_drive = PathBuf::from(format!("{}.retired", replaced_drive.display())); std::fs::rename(&replaced_drive, &retired_drive)?; @@ -138,6 +295,11 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res .await?; dist.cluster.start_node(replaced_node).await?; + let target_pid_after = dist.cluster.nodes[replaced_node] + .process + .as_ref() + .ok_or("target process is absent after restart")? + .id(); let heal_body = r#"{"recursive":true,"dryRun":false,"remove":false,"recreate":true,"scanMode":2,"updateParity":false,"nolock":false}"#; let heal_url = format!("{}/rustfs/admin/v3/heal/{bucket}?forceStart=true", dist.cluster.nodes[0].url); @@ -167,6 +329,7 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res .chain(std::iter::once((outage_key.to_string(), outage_body.clone()))) .collect::>(); let expected_keys = inventory.keys().cloned().collect::>(); + let mut node_listings = Vec::new(); for node_index in 0..dist.cluster.nodes.len() { let client = dist.client(node_index)?; assert_inventory(&client, &bucket, &inventory).await?; @@ -177,6 +340,24 @@ async fn three_node_four_drive_ec8_4_root_heal_rebuilds_replaced_drive_after_res .filter_map(|object| object.key().map(str::to_owned)) .collect::>(); assert_eq!(observed, expected_keys, "node {node_index} listing diverged after EC8+4 heal"); + let mut observed = observed.into_iter().collect::>(); + observed.sort(); + node_listings.push(observed); + } + if let Some(context) = evidence_context { + write_scanner_heal_evidence( + context, + &dist, + &bucket, + &expected, + outage_key, + &outage_body, + &replaced_drive, + target_pid_before, + target_pid_after, + node_listings, + ) + .await?; } Ok(()) From 4a2b15cb82c1307aa52cd83630c6afed329a717e Mon Sep 17 00:00:00 2001 From: houseme Date: Tue, 8 Sep 2026 18:53:43 +0800 Subject: [PATCH 06/12] fix(heal): harden MRF replay boundaries (#7483) Reject journal records with unknown version-presence flags even when their CRC is valid, so rollback/future payloads cannot be accepted as known records. Gate committed checkpoint cleanup by the writer owner captured from the replay source, preserving retained manifests from other owners inside the same sequence window. Co-authored-by: zhi22915 --- crates/heal/src/heal/mrf_queue.rs | 52 +++++++++++++++----- crates/heal/src/heal/mrf_queue/snapshot.rs | 56 ++++++++++++++++++++-- 2 files changed, 90 insertions(+), 18 deletions(-) diff --git a/crates/heal/src/heal/mrf_queue.rs b/crates/heal/src/heal/mrf_queue.rs index 48ae146ce..77d52720b 100644 --- a/crates/heal/src/heal/mrf_queue.rs +++ b/crates/heal/src/heal/mrf_queue.rs @@ -294,7 +294,11 @@ fn decode_one(data: &[u8]) -> Option<(MrfIntent, usize)> { }; let attempts = data[3]; let enqueued_at_ms = u64::from_le_bytes(data[4..12].try_into().ok()?); - let has_version = data[12] != 0; + let has_version = match data[12] { + 0 => false, + 1 => true, + _ => return None, + }; let mut cursor = MRF_RECORD_FIXED_HEAD; let version_id = if has_version { if data.len() < cursor + 16 { @@ -718,7 +722,7 @@ fn replay_must_retain_journal( #[derive(Clone, Copy)] enum ReplayCleanup { Legacy, - Committed { sequence: u64 }, + Committed { owner: Uuid, sequence: u64 }, } struct ReplaySource { @@ -731,6 +735,7 @@ async fn read_replay_source(max_bytes: usize) -> Result, sn return Ok(Some(ReplaySource { data: committed.payload().to_vec(), cleanup: ReplayCleanup::Committed { + owner: committed.owner(), sequence: committed.sequence(), }, })); @@ -754,18 +759,20 @@ async fn read_replay_source(max_bytes: usize) -> Result, sn async fn delete_replay_source(cleanup: ReplayCleanup, max_bytes: usize) -> bool { let committed_deleted = match cleanup { ReplayCleanup::Legacy => true, - ReplayCleanup::Committed { sequence } => match snapshot::delete_committed_snapshots_through(sequence, max_bytes).await { - Ok(deleted) => deleted, - Err(err) => { - tracing::warn!( - target: "rustfs::heal::mrf", - error = %err, - sequence, - "MRF committed replay checkpoint cleanup failed" - ); - false + ReplayCleanup::Committed { owner, sequence } => { + match snapshot::delete_committed_snapshots_through(owner, sequence, max_bytes).await { + Ok(deleted) => deleted, + Err(err) => { + tracing::warn!( + target: "rustfs::heal::mrf", + error = %err, + sequence, + "MRF committed replay checkpoint cleanup failed" + ); + false + } } - }, + } }; committed_deleted && delete_journals().await } @@ -1344,6 +1351,25 @@ mod tests { assert_eq!(truncated, corrupt.len()); } + #[test] + fn journal_rejects_unknown_version_presence_flag_even_with_valid_crc() { + let mut versioned = intent("rollback-bucket", "object", 0); + versioned.version_id = Some([9; 16]); + let mut buf = Vec::new(); + assert!(encode_intent(&versioned, &mut buf)); + + buf[12] = 2; + let crc_offset = buf.len() - 4; + let mut hasher = crc_fast::Digest::new(crc_fast::CrcAlgorithm::Crc32IsoHdlc); + hasher.update(&buf[..crc_offset]); + let checksum = u32::try_from(hasher.finalize()).expect("CRC32 fits"); + buf[crc_offset..].copy_from_slice(&checksum.to_le_bytes()); + + let (decoded, truncated) = decode_journal(&buf); + assert!(decoded.is_empty(), "unknown boolean encodings are not rollback-compatible payloads"); + assert_eq!(truncated, buf.len()); + } + #[test] fn heal_request_mapping_follows_priority_matrix() { let decode = build_heal_request(&intent("b", "o", 0)); diff --git a/crates/heal/src/heal/mrf_queue/snapshot.rs b/crates/heal/src/heal/mrf_queue/snapshot.rs index 0921c043f..d4cafac10 100644 --- a/crates/heal/src/heal/mrf_queue/snapshot.rs +++ b/crates/heal/src/heal/mrf_queue/snapshot.rs @@ -555,20 +555,25 @@ pub async fn inspect_local_committed_snapshot(max_bytes: usize) -> Result