From 84418df1046e7ef085fb7a981340d2403f65e737 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=A9=AC=E7=99=BB=E5=B1=B1?= Date: Thu, 30 Jul 2026 09:23:55 +0800 Subject: [PATCH] fix(ecstore): reject poisoned format heal candidates --- crates/ecstore/src/core/sets.rs | 185 +++++++++++++++++++----- crates/ecstore/src/set_disk/mod.rs | 3 +- crates/ecstore/src/set_disk/ops/heal.rs | 22 +-- crates/ecstore/src/store/heal.rs | 71 ++++++++- 4 files changed, 230 insertions(+), 51 deletions(-) diff --git a/crates/ecstore/src/core/sets.rs b/crates/ecstore/src/core/sets.rs index 65cef6e7f..7b4ba2a60 100644 --- a/crates/ecstore/src/core/sets.rs +++ b/crates/ecstore/src/core/sets.rs @@ -37,7 +37,9 @@ use crate::{ runtime::instance::{InstanceContext, bootstrap_ctx}, runtime::sources as runtime_sources, set_disk::{PreparedGetObjectMetadata, SetDisks}, - store::init_format::{check_format_erasure_values, get_format_erasure_in_quorum, load_format_erasure_all, save_format_file}, + store::init_format::{ + check_format_erasure_values, load_format_erasure_all, save_format_file, select_format_erasure_in_quorum, + }, }; use futures::{ future::join_all, @@ -947,7 +949,7 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets { #[tracing::instrument(skip(self))] async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { - let (disks, _) = init_storage_disks_with_errors( + let (disks, init_errs) = init_storage_disks_with_errors( &self.endpoints.endpoints, &DiskOption { cleanup: false, @@ -955,16 +957,36 @@ impl crate::storage_api_contracts::heal::HealOperations for Sets { }, ) .await; - let (formats, errs) = load_format_erasure_all(&disks, true).await; + let (formats, mut errs) = load_format_erasure_all(&disks, true).await; + for (err, init_err) in errs.iter_mut().zip(init_errs) { + if init_err.is_some() { + *err = init_err; + } + } + if errs.iter().any(|err| { + matches!( + err, + Some(DiskError::InconsistentDisk | DiskError::CorruptedFormat | DiskError::CorruptedBackend) + ) + }) { + return Ok((HealResultItem::default(), Some(StorageError::CorruptedFormat))); + } if let Err(err) = check_format_erasure_values(&formats, self.set_drive_count) { info!("failed to check formats erasure values: {}", err); return Ok((HealResultItem::default(), Some(err))); } - let ref_format = match get_format_erasure_in_quorum(&formats, 0) { - Ok(format) if format.shared_identity() == self.format.shared_identity() => format, + let (ref_format, quorum_members) = match select_format_erasure_in_quorum(&formats, 0) { + Ok((format, members)) if format.shared_identity() == self.format.shared_identity() => (format, members), Ok(_) => return Ok((HealResultItem::default(), Some(StorageError::CorruptedFormat))), Err(err) => return Ok((HealResultItem::default(), Some(err))), }; + if formats + .iter() + .zip(quorum_members) + .any(|(format, member)| format.is_some() && !member) + { + return Ok((HealResultItem::default(), Some(StorageError::CorruptedFormat))); + } let mut res = HealResultItem { heal_item_type: HealItemType::Metadata.to_string(), detail: "disk-format".to_string(), @@ -1385,34 +1407,8 @@ mod tests { #[tokio::test] async fn format_heal_rejects_foreign_majorities_at_set_and_pool_scopes() { - let (_temp_dirs, canonical_format, sets) = setup_heal_format_sets(2, true).await; - let endpoints = sets.endpoints.endpoints.as_ref().clone(); - let mut disks = Vec::with_capacity(endpoints.len()); - for endpoint in &endpoints { - disks.push(Some( - new_disk( - endpoint, - &DiskOption { - cleanup: false, - health_check: false, - }, - ) - .await - .expect("fresh set-level disk handle should open"), - )); - } - let set_disks = SetDisks::new( - "test-owner".to_string(), - Arc::new(RwLock::new(disks)), - endpoints.len(), - 1, - 0, - 0, - endpoints, - canonical_format, - Vec::new(), - ) - .await; + let (_temp_dirs, _canonical_format, sets) = setup_heal_format_sets(2, true).await; + let set_disks = set_level_heal_view(&sets).await; let (_, set_err) = set_disks .heal_format(false) @@ -1433,6 +1429,71 @@ mod tests { ); } + #[tokio::test] + async fn pool_format_heal_rejects_a_wrong_slot_minority() { + let (_temp_dirs, canonical_format, sets) = setup_heal_format_sets(3, false).await; + let mut poisoned_format = canonical_format.clone(); + poisoned_format.erasure.this = canonical_format.erasure.sets[0][0]; + replace_heal_test_format(&sets, 2, &poisoned_format).await; + let probe_err = new_disk( + &sets.endpoints.endpoints.as_ref()[2], + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect_err("a wrong-slot local format must fail disk initialization"); + assert_eq!(probe_err, DiskError::InconsistentDisk); + + let (_, pool_err) = sets + .heal_format(false) + .await + .expect("pool format heal should report a typed slot mismatch"); + assert!( + matches!(pool_err, Some(StorageError::CorruptedFormat)), + "a wrong-slot minority must not be reported as no-heal-required: {pool_err:?}" + ); + assert_eq!( + read_heal_test_format(&sets, 2).await, + poisoned_format, + "format heal must not overwrite a wrong-slot disk" + ); + } + + #[tokio::test] + async fn format_heal_rejects_a_foreign_minority_at_set_and_pool_scopes() { + let (_temp_dirs, canonical_format, sets) = setup_heal_format_sets(3, false).await; + let mut poisoned_format = canonical_format.clone(); + poisoned_format.id = Uuid::new_v4(); + poisoned_format.erasure.this = poisoned_format.erasure.sets[0][2]; + replace_heal_test_format(&sets, 2, &poisoned_format).await; + let set_disks = set_level_heal_view(&sets).await; + + let (_, set_err) = set_disks + .heal_format(false) + .await + .expect("set format heal should report a typed identity mismatch"); + assert!( + matches!(set_err, Some(StorageError::CorruptedFormat)), + "a foreign minority must not be reported as no-heal-required: {set_err:?}" + ); + + let (_, pool_err) = sets + .heal_format(false) + .await + .expect("pool format heal should report a typed identity mismatch"); + assert!( + matches!(pool_err, Some(StorageError::CorruptedFormat)), + "a foreign minority must not be reported as no-heal-required: {pool_err:?}" + ); + assert_eq!( + read_heal_test_format(&sets, 2).await, + poisoned_format, + "format heal must not overwrite a foreign disk" + ); + } + #[tokio::test(flavor = "multi_thread")] #[serial] async fn list_multipart_uploads_merges_all_sets_without_pagination_loss() { @@ -1676,8 +1737,8 @@ mod tests { // formatting the first `num_formatted` of them against a shared reference // format and leaving the rest unformatted. Returns the live TempDir handles // (must be kept alive), the reference format, and the assembled `Sets`. - // `disk_set` is intentionally empty: these tests only drive `heal_format` - // with `dry_run == true`, which never touches `disk_set`. + // `disk_set` is intentionally empty: these tests only exercise paths that + // return before pool-level healing delegates into a set. async fn setup_heal_format_sets(num_formatted: usize, foreign_identity: bool) -> (Vec, FormatV3, Sets) { const SET_DRIVE_COUNT: usize = 3; let ref_format = FormatV3::new(1, SET_DRIVE_COUNT); @@ -1741,6 +1802,60 @@ mod tests { (dirs, ref_format, sets) } + async fn set_level_heal_view(sets: &Sets) -> Arc { + let endpoints = sets.endpoints.endpoints.as_ref().clone(); + let mut disks = Vec::with_capacity(endpoints.len()); + for endpoint in &endpoints { + disks.push(Some( + new_disk( + endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("fresh set-level disk handle should open"), + )); + } + + SetDisks::new( + "test-owner".to_string(), + Arc::new(RwLock::new(disks)), + endpoints.len(), + 1, + 0, + 0, + endpoints, + sets.format.clone(), + Vec::new(), + ) + .await + } + + async fn replace_heal_test_format(sets: &Sets, disk_index: usize, format: &FormatV3) { + let disk = new_disk( + &sets.endpoints.endpoints.as_ref()[disk_index], + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("heal test disk should open"); + save_format_file(&Some(disk.clone()), &Some(format.clone())) + .await + .expect("poisoned test format should be written"); + } + + async fn read_heal_test_format(sets: &Sets, disk_index: usize) -> FormatV3 { + let path = std::path::Path::new(&sets.endpoints.endpoints.as_ref()[disk_index].get_file_path()) + .join(crate::disk::RUSTFS_META_BUCKET) + .join(crate::disk::FORMAT_CONFIG_FILE); + let data = tokio::fs::read(path).await.expect("test format should be readable"); + FormatV3::try_from(data.as_slice()).expect("test format should parse") + } + // Regression for #956 (NoHealRequired path): with every disk already // formatted, `heal_format` reports exactly one drive record per disk // (N = set_count * set_drive_count), each carrying a real endpoint. Before diff --git a/crates/ecstore/src/set_disk/mod.rs b/crates/ecstore/src/set_disk/mod.rs index 538d9253b..d356108d3 100644 --- a/crates/ecstore/src/set_disk/mod.rs +++ b/crates/ecstore/src/set_disk/mod.rs @@ -111,7 +111,8 @@ use crate::{ // event::name::EventName, services::event_notification::{EventArgs, send_event}, store::init_format::{ - format_disk_id_matches_slot, get_format_erasure_in_quorum, load_format_erasure, load_format_erasure_all, save_format_file, + formats_match_reference_slots, get_format_erasure_in_quorum, load_format_erasure, load_format_erasure_all, + save_format_file, }, }; use bytes::Bytes; diff --git a/crates/ecstore/src/set_disk/ops/heal.rs b/crates/ecstore/src/set_disk/ops/heal.rs index f7048c4dd..91e8570d0 100644 --- a/crates/ecstore/src/set_disk/ops/heal.rs +++ b/crates/ecstore/src/set_disk/ops/heal.rs @@ -1373,6 +1373,14 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks { async fn heal_format(&self, dry_run: bool) -> Result<(HealResultItem, Option)> { let disks = self.disks.read().await.clone(); let (formats, errs) = load_format_erasure_all(&disks, true).await; + if errs.iter().any(|err| { + matches!( + err, + Some(DiskError::InconsistentDisk | DiskError::CorruptedFormat | DiskError::CorruptedBackend) + ) + }) { + return Ok((HealResultItem::default(), Some(StorageError::CorruptedFormat))); + } let slot_offset = self .set_index .checked_mul(self.set_drive_count) @@ -1382,14 +1390,7 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks { Ok(_) => return Ok((HealResultItem::default(), Some(StorageError::CorruptedFormat))), Err(err) => { let can_use_cached_layout = count_errs(&errs, &DiskError::UnformattedDisk) > 0 - && formats.iter().enumerate().all(|(index, format)| { - format.as_ref().is_none_or(|format| { - self.format.shared_identity() == format.shared_identity() - && slot_offset - .checked_add(index) - .is_some_and(|slot| format_disk_id_matches_slot(format, slot)) - }) - }) + && formats_match_reference_slots(&formats, &self.format, slot_offset) && errs .iter() .all(|err| err.is_none() || matches!(err, Some(DiskError::UnformattedDisk))); @@ -1400,6 +1401,9 @@ impl crate::storage_api_contracts::heal::HealOperations for SetDisks { } } }; + if !formats_match_reference_slots(&formats, &ref_format, slot_offset) { + return Ok((HealResultItem::default(), Some(StorageError::CorruptedFormat))); + } let endpoints = crate::layout::endpoints::Endpoints::from(self.set_endpoints.clone()); let before_drives = crate::layout::set_heal::formats_to_drives_info(&endpoints, &formats, &errs); @@ -1911,7 +1915,7 @@ mod heal_result_report_tests { .await .expect("format heal should report the quorum failure in its result"); - assert!(matches!(heal_err, Some(Error::ErasureReadQuorum))); + assert!(matches!(heal_err, Some(Error::CorruptedFormat))); let unformatted = load_format_erasure(disks[1].as_ref().expect("second disk should be online"), true) .await .expect_err("a rejected fallback must not format the missing slot"); diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index bedfa6fd8..fe99c1291 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -30,6 +30,7 @@ impl ECStore { }; let mut count_no_heal = 0; + let mut first_error = None; for pool in self.pools.iter() { let (mut result, err) = pool.heal_format(dry_run).await?; if let Some(err) = err { @@ -37,7 +38,9 @@ impl ECStore { StorageError::NoHealRequired => { count_no_heal += 1; } - err => return Ok((result, Some(err))), + err => { + first_error.get_or_insert(err); + } } } r.disk_count += result.disk_count; @@ -45,6 +48,9 @@ impl ECStore { r.before.drives.append(&mut result.before.drives); r.after.drives.append(&mut result.after.drives); } + if let Some(err) = first_error { + return Ok((r, Some(err))); + } if count_no_heal == self.pools.len() { info!( event = EVENT_HEAL_FORMAT_COMPLETED, @@ -169,10 +175,10 @@ mod tests { use super::*; use crate::disk::{DiskOption, format::FormatV3, new_disk}; use crate::layout::endpoints::{Endpoints, PoolEndpoints}; - use crate::store::init_format::save_format_file; + use crate::store::init_format::{load_format_erasure, save_format_file}; #[tokio::test] - async fn handle_heal_format_propagates_a_foreign_format_majority() { + async fn handle_heal_format_continues_after_a_pool_error() { let canonical_format = FormatV3::new(1, 3); let mut foreign_format = canonical_format.clone(); foreign_format.id = Uuid::new_v4(); @@ -214,14 +220,62 @@ mod tests { cmd_line: "foreign-format-majority-test".to_string(), platform: "test".to_string(), }; - let endpoint_pools = EndpointServerPools::from(vec![pool_endpoints.clone()]); let pool = Sets::new(disks, &pool_endpoints, &canonical_format, 0, 1) .await .expect("test pool should build around the cached canonical format"); + + let mut recoverable_format = FormatV3::new(1, 3); + recoverable_format.id = canonical_format.id; + let mut recoverable_temp_dirs = Vec::new(); + let mut recoverable_endpoints = Vec::new(); + let mut recoverable_disks = Vec::new(); + let mut unformatted_disk = None; + for disk_index in 0..3 { + let temp_dir = tempfile::tempdir().expect("temporary disk root should be created"); + let mut endpoint = Endpoint::try_from(temp_dir.path().to_str().expect("temporary path should be UTF-8")) + .expect("temporary endpoint should parse"); + endpoint.set_pool_index(1); + endpoint.set_set_index(0); + endpoint.set_disk_index(disk_index); + let disk = new_disk( + &endpoint, + &DiskOption { + cleanup: false, + health_check: false, + }, + ) + .await + .expect("temporary disk should open"); + if disk_index < 2 { + let mut disk_format = recoverable_format.clone(); + disk_format.erasure.this = recoverable_format.erasure.sets[0][disk_index]; + save_format_file(&Some(disk.clone()), &Some(disk_format)) + .await + .expect("recoverable format should be written"); + } else { + unformatted_disk = Some(disk.clone()); + } + recoverable_temp_dirs.push(temp_dir); + recoverable_endpoints.push(endpoint); + recoverable_disks.push(Some(disk)); + } + let recoverable_pool_endpoints = PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 3, + endpoints: Endpoints::from(recoverable_endpoints), + cmd_line: "recoverable-format-test".to_string(), + platform: "test".to_string(), + }; + let recoverable_pool = Sets::new(recoverable_disks, &recoverable_pool_endpoints, &recoverable_format, 1, 1) + .await + .expect("recoverable test pool should build"); + + let endpoint_pools = EndpointServerPools::from(vec![pool_endpoints.clone(), recoverable_pool_endpoints.clone()]); let store = ECStore { id: canonical_format.id, disk_map: HashMap::new(), - pools: vec![pool], + pools: vec![pool, recoverable_pool], peer_sys: S3PeerSys::new(&endpoint_pools), pool_meta: RwLock::new(PoolMeta::default()), rebalance_meta: RwLock::new(None), @@ -231,7 +285,7 @@ mod tests { ctx: crate::runtime::instance::bootstrap_ctx(), }; - let (_, err) = store + let (result, err) = store .handle_heal_format(false) .await .expect("format heal should return the typed pool error"); @@ -239,5 +293,10 @@ mod tests { matches!(err, Some(StorageError::CorruptedFormat)), "foreign format majority must not be downgraded to a successful heal: {err:?}" ); + assert_eq!(result.disk_count, 3, "the recoverable pool should still be inspected"); + let healed = load_format_erasure(&unformatted_disk.expect("the unformatted disk handle should be retained"), true) + .await + .expect("the later pool should be healed despite the first pool error"); + assert_eq!(healed.erasure.this, recoverable_format.erasure.sets[0][2]); } }