From c03d3cdd59df07d0db7bd6ec15c207f5d678c073 Mon Sep 17 00:00:00 2001 From: overtrue Date: Tue, 8 Sep 2026 17:23:53 +0800 Subject: [PATCH] fix: close replacement and protocol validation gaps --- .../src/heal_erasure_disk_rebuild_test.rs | 8 +- .../e2e_test/src/protocols/sftp_compliance.rs | 7 +- crates/e2e_test/src/protocols/sftp_core.rs | 17 +-- crates/ecstore/src/store/heal.rs | 87 ++++++++++++ crates/heal/src/heal/erasure_healer.rs | 125 ++++++++++++++++++ crates/heal/src/heal/storage.rs | 20 +++ 6 files changed, 250 insertions(+), 14 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 a86f82473..06f443276 100644 --- a/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs +++ b/crates/e2e_test/src/heal_erasure_disk_rebuild_test.rs @@ -1499,7 +1499,11 @@ mod tests { "Restored target endpoint forwarding" ); } else { - if scenario == InterruptionScenario::BackgroundTargetRestart { + let graceful_restart = matches!( + scenario, + InterruptionScenario::BackgroundTargetRestart | InterruptionScenario::BackgroundTargetRestartEc84 + ); + if graceful_restart { cluster.stop_node_gracefully(interruption_node).await?; } else { cluster.stop_node(interruption_node)?; @@ -1521,7 +1525,7 @@ 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 = !graceful_restart; assert!( marker_exists == expected_marker, "background restart/crash lane observed unexpected unclean-shutdown marker state" diff --git a/crates/e2e_test/src/protocols/sftp_compliance.rs b/crates/e2e_test/src/protocols/sftp_compliance.rs index be8af0e87..eea3195ff 100644 --- a/crates/e2e_test/src/protocols/sftp_compliance.rs +++ b/crates/e2e_test/src/protocols/sftp_compliance.rs @@ -79,6 +79,11 @@ pub async fn test_sftp_compliance_suite() -> Result<()> { .await .map_err(|e| anyhow!("{}", e))?; + // Protocol listeners can accept connections before IAM is initialized. + // A signed S3 request establishes readiness before the first SFTP login. + let s3 = build_test_s3_client(&format!("http://{COMPLIANCE_RW_S3_ADDRESS}")); + wait_for_s3_ready(&s3, 30).await?; + let (session, sftp) = connect_sftp_to(COMPLIANCE_RW_SFTP_ADDRESS).await?; cmptst_01::run_medium_binary_round_trip(&sftp).await?; @@ -101,8 +106,6 @@ pub async fn test_sftp_compliance_suite() -> Result<()> { // reach the finalised object as x-amz-meta-* user metadata // through the CreateMultipartUpload input field. The S3 client // connects to the same rustfs process this suite already drives. - let s3 = build_test_s3_client(&format!("http://{COMPLIANCE_RW_S3_ADDRESS}")); - wait_for_s3_ready(&s3, 30).await?; cmptst_34::run_open_attrs_round_trip_multipart(&sftp, &s3).await?; drop(sftp); diff --git a/crates/e2e_test/src/protocols/sftp_core.rs b/crates/e2e_test/src/protocols/sftp_core.rs index a32925fa4..b43b3ae7e 100644 --- a/crates/e2e_test/src/protocols/sftp_core.rs +++ b/crates/e2e_test/src/protocols/sftp_core.rs @@ -168,6 +168,10 @@ pub async fn test_sftp_core_operations() -> Result<()> { .await .map_err(|e| anyhow!("{}", e))?; + // Protocol listeners can accept connections before IAM is initialized. + let s3 = build_test_s3_client(S3_ENDPOINT); + wait_for_s3_ready(&s3, S3_READY_ATTEMPTS).await?; + let (session, sftp) = connect_sftp().await?; // --- 1. Subsystem canary: SFTP session reachable after password auth --- @@ -348,16 +352,6 @@ pub async fn test_sftp_core_operations() -> Result<()> { let _ = bad_session.disconnect(russh::Disconnect::ByApplication, "", "en").await; info!("PASS: bad-password authentication rejected"); - // --- Cross-protocol setup: aws-sdk-s3 client against the same server --- - // The rustfs binary spawned for this suite serves both SFTP on port - // 9022 and S3 on port 9000. The S3 stack may need a moment to finish - // initialising after TCP is listening, so list_buckets is polled - // until it succeeds before any cross-protocol assertion runs. - info!("Testing SFTP: prepare aws-sdk-s3 client and wait for S3 readiness"); - let s3 = build_test_s3_client(S3_ENDPOINT); - wait_for_s3_ready(&s3, S3_READY_ATTEMPTS).await?; - info!("PASS: S3 endpoint reachable from cross-protocol client"); - // --- SFTP write, S3 read: SHA256 round-trip --- // SFTP creates the object, then assert_cross_protocol_sha_match // fetches it via both S3 GetObject and SFTP READ and compares @@ -522,6 +516,9 @@ pub async fn test_sftp_idle_timeout_disconnects() -> Result<()> { .await .map_err(|e| anyhow!("{}", e))?; + let s3 = build_test_s3_client(&format!("http://{IDLE_S3_ADDRESS}")); + wait_for_s3_ready(&s3, S3_READY_ATTEMPTS).await?; + let (session, sftp) = connect_sftp_to(IDLE_SFTP_ADDRESS).await?; // Confirm the session is live before the wait so a failure in the diff --git a/crates/ecstore/src/store/heal.rs b/crates/ecstore/src/store/heal.rs index 8144e7eea..edd5bb673 100644 --- a/crates/ecstore/src/store/heal.rs +++ b/crates/ecstore/src/store/heal.rs @@ -378,6 +378,26 @@ impl ECStore { Ok(result) } + /// Whether this replacement set owns the pool's metadata replica. + /// + /// Pool metadata follows normal object placement within each pool. A valid + /// non-owner set has no replica to repair; missing metadata on the owner + /// set still requires healing and target-specific readback. + pub fn replacement_pool_metadata_applies(&self, pool_index: usize, set_index: usize) -> Result { + let pool = self + .pools + .get(pool_index) + .ok_or_else(|| invalid_heal_pool_index(pool_index, self.pools.len()))?; + let selected = pool.get_disks_for_heal_object( + POOL_META_NAME, + &HealOpts { + set: Some(set_index), + ..Default::default() + }, + )?; + Ok(Arc::ptr_eq(&selected, &pool.get_disks_by_key(POOL_META_NAME))) + } + #[instrument(skip(self, targets), fields(pool_index, set_index, target_count = targets.len()))] pub async fn replacement_targets_have_version( &self, @@ -829,6 +849,73 @@ mod tests { } } + #[tokio::test] + async fn replacement_pool_metadata_applies_to_the_written_replica_in_each_pool() { + let mut store = minimal_heal_store().await; + for pool_index in 0..store.pools.len() { + assert!( + store + .replacement_pool_metadata_applies(pool_index, 0) + .expect("a valid single-set pool should have a metadata owner") + ); + } + store.ctx = Arc::new(InstanceContext::new()); + for algorithm in [ + crate::disk::format::DistributionAlgoVersion::V1, + crate::disk::format::DistributionAlgoVersion::V2, + crate::disk::format::DistributionAlgoVersion::V3, + ] { + let mut temp_dirs = Vec::new(); + for pool_index in 0..store.pools.len() { + let (dirs, mut pool) = + crate::core::sets::make_local_two_set_sets_for_pool_with_ctx(Arc::clone(&store.ctx), pool_index).await; + temp_dirs.extend(dirs); + Arc::get_mut(&mut pool) + .expect("fixture pool should have one owner") + .distribution_algo = algorithm.clone(); + store.pools[pool_index] = pool; + } + for (pool_index, pool) in store.pools.iter().enumerate() { + let mut required_sets = 0; + for set_index in 0..pool.disk_set.len() { + required_sets += usize::from( + store + .replacement_pool_metadata_applies(pool_index, set_index) + .expect("valid replacement topology should be classified before metadata exists"), + ); + } + assert_eq!(required_sets, 1, "missing metadata cannot exempt the owner set"); + save_config(pool.clone(), POOL_META_NAME, b"pool metadata placement".to_vec()) + .await + .expect("normal config writes should persist one metadata replica per pool"); + for (set_index, set) in pool.disk_set.iter().enumerate() { + let applies = store + .replacement_pool_metadata_applies(pool_index, set_index) + .expect("valid replacement topology should be classified"); + let disks = set.disks.read().await.clone(); + for disk in disks.iter().flatten() { + let replica = disk.read_xl(RUSTFS_META_BUCKET, POOL_META_NAME, false).await; + if applies { + replica.expect("the metadata owner must match actual persisted shards"); + } else { + assert!( + matches!(replica, Err(crate::disk::error::DiskError::FileNotFound)), + "non-owner sets must have no persisted metadata shard; observed error: {:?}", + replica.as_ref().err() + ); + } + } + } + assert!( + store + .replacement_pool_metadata_applies(pool_index, pool.disk_set.len()) + .is_err() + ); + } + } + assert!(store.replacement_pool_metadata_applies(store.pools.len(), 0).is_err()); + } + async fn remove_pool_meta_shard(store: &ECStore, pool_idx: usize) -> DiskStore { let target_set = store.pools[pool_idx].get_disks_by_key(POOL_META_NAME); let missing_disk = target_set.disks.read().await[0] diff --git a/crates/heal/src/heal/erasure_healer.rs b/crates/heal/src/heal/erasure_healer.rs index dbd62e422..df72f281c 100644 --- a/crates/heal/src/heal/erasure_healer.rs +++ b/crates/heal/src/heal/erasure_healer.rs @@ -957,6 +957,10 @@ impl ErasureSetHealer { }); } + if !self.storage.replacement_pool_metadata_applies(&self.heal_opts).await? { + return Ok(()); + } + let object_key = format!("{RUSTFS_META_BUCKET}/{POOL_META_NAME}"); let checkpoint_key = compose_key(&object_key, None); let checkpoint = checkpoint_manager.get_checkpoint().await; @@ -2029,6 +2033,8 @@ mod resume_loop_tests { #[derive(Clone)] enum HealOutcome { Ok, + /// The object has no metadata on any disk in the selected set. + FileNotFound, /// The version vanished before heal ran (deleted mid-heal). VersionNotFound, /// A transient infrastructure condition (offline disk / unmet quorum): @@ -2054,6 +2060,8 @@ mod resume_loop_tests { /// Target-specific physical readback evidence per `compose_key`; the /// fake models a healthy backend unless a test explicitly revokes it. replacement_commit_evidence: Mutex>, + pool_metadata_not_applicable: AtomicBool, + fail_pool_metadata_scope: AtomicBool, lifecycle_expired: Mutex>, /// every heal_object call recorded as (name, version_id) heal_calls: Mutex)>>, @@ -2155,6 +2163,7 @@ mod resume_loop_tests { let outcome = self.outcomes.lock().unwrap().get(&key).cloned().unwrap_or(HealOutcome::Ok); match outcome { HealOutcome::Ok => Ok((self.results.lock().unwrap().get(&key).cloned().unwrap_or_default(), None)), + HealOutcome::FileNotFound => Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::FileNotFound)))), HealOutcome::VersionNotFound => { Ok((HealResultItem::default(), Some(Error::Storage(EcstoreError::FileVersionNotFound)))) } @@ -2168,6 +2177,17 @@ mod resume_loop_tests { async fn heal_format(&self, _dry: bool) -> Result<(HealResultItem, Option)> { Ok((HealResultItem::default(), None)) } + async fn replacement_pool_metadata_applies(&self, opts: &HealOpts) -> Result { + if self.fail_pool_metadata_scope.load(Ordering::SeqCst) { + return Err(Error::other("injected pool metadata scope failure")); + } + if self.pool_metadata_not_applicable.load(Ordering::SeqCst) { + assert_eq!(opts.pool, Some(0)); + assert_eq!(opts.set, Some(1)); + return Ok(false); + } + Ok(true) + } async fn replacement_targets_have_version( &self, _bucket: &str, @@ -2713,6 +2733,111 @@ mod resume_loop_tests { drop(checkpoint); } + #[tokio::test] + async fn replacement_pool_metadata_non_owner_completes_but_missing_owner_retries() { + for owns_pool_metadata in [false, true] { + let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; + let replacement_task_id = ResumeUtils::generate_task_id(); + let set_index = usize::from(!owns_pool_metadata); + let set_disk_id = format!("pool_0_set_{set_index}"); + ResumeManager::new_replacement_intent( + env.healer.disk.clone(), + replacement_task_id.clone(), + set_disk_id.clone(), + vec!["b".to_string()], + vec!["replacement-a".to_string()], + vec![crate::heal::resume::ReplacementTargetIdentity { + endpoint: "replacement-a".to_string(), + canonical_path: "/mnt/replacement-a".to_string(), + physical_device_ids: vec!["device-a".to_string()], + filesystem_identity: "1:2:3".to_string(), + }], + ) + .await + .expect("replacement intent should persist"); + env.storage + .pool_metadata_not_applicable + .store(!owns_pool_metadata, Ordering::SeqCst); + env.storage.set_outcome(POOL_META_NAME, None, HealOutcome::FileNotFound); + let healer = ErasureSetHealer::new( + env.storage.clone(), + Arc::new(RwLock::new(HealProgress::new())), + CancellationToken::new(), + env.healer.disk.clone(), + HealOpts { + pool: Some(0), + set: Some(set_index), + ..Default::default() + }, + HealRequestSource::AutoHeal, + ) + .with_replacement_targets(vec!["replacement-a".to_string()], Some(replacement_task_id.clone())); + + let result = healer.heal_erasure_set(&["b".to_string()], &set_disk_id).await; + let state = ResumeManager::load_replacement_intent(env.healer.disk.clone(), &replacement_task_id) + .await + .expect("replacement state must remain until marker cleanup") + .get_state() + .await; + if owns_pool_metadata { + let error = result.expect_err("missing metadata in the owner set must keep replacement incomplete"); + assert!(error.to_string().contains("Replacement erasure set heal incomplete")); + assert!(!state.completed); + assert_eq!(state.replacement_phase, crate::heal::resume::ReplacementPhase::Intent); + assert_eq!(state.retry_count, 1); + assert_eq!(env.storage.calls(), vec![(POOL_META_NAME.to_string(), None)]); + } else { + result.expect("a non-owner set must complete without a pool metadata replica"); + assert!(state.completed); + assert_eq!(state.replacement_phase, crate::heal::resume::ReplacementPhase::Verified); + assert_eq!(state.retry_count, 0); + assert!(env.storage.calls().is_empty(), "non-owner sets must not attempt pool metadata repair"); + } + } + } + + #[tokio::test] + async fn replacement_pool_metadata_unknown_scope_cannot_complete() { + 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 { + pool: Some(0), + set: Some(0), + ..Default::default() + }, + HealRequestSource::AutoHeal, + ) + .with_replacement_targets(vec!["replacement-a".to_string()], Some("generation-a".to_string())); + env.storage.fail_pool_metadata_scope.store(true, Ordering::SeqCst); + env.storage + .set_result(POOL_META_NAME, None, replacement_target_ok_result("replacement-a", POOL_META_NAME)); + let mut processed_objects = 0; + let mut successful_objects = 0; + let mut failed_objects = 0; + let mut skipped_objects = 0; + let error = healer + .heal_replacement_pool_metadata( + "pool_0_set_0", + &mut super::ErasureSetPassCounters { + processed_objects: &mut processed_objects, + successful_objects: &mut successful_objects, + failed_objects: &mut failed_objects, + skipped_objects: &mut skipped_objects, + }, + &env.resume, + &env.checkpoint, + ) + .await + .expect_err("unknown metadata placement must keep replacement incomplete"); + assert!(error.to_string().contains("injected pool metadata scope failure")); + assert!(env.storage.calls().is_empty()); + assert_eq!((processed_objects, successful_objects, failed_objects, skipped_objects), (0, 0, 0, 0)); + } + #[tokio::test] async fn replacement_pool_metadata_readback_failure_schedules_retry() { let env = make_env_with_targets(vec!["replacement-a".to_string()]).await; diff --git a/crates/heal/src/heal/storage.rs b/crates/heal/src/heal/storage.rs index 1dc6fdddc..63844c247 100644 --- a/crates/heal/src/heal/storage.rs +++ b/crates/heal/src/heal/storage.rs @@ -436,6 +436,14 @@ pub trait HealStorageAPI: Send + Sync { Err(Error::other("target-scoped replacement format is unsupported")) } + /// Whether the selected replacement set owns the pool metadata replica. + /// + /// Only a topology-aware backend may exempt a valid non-owner set. The + /// conservative default requires the existing repair and readback checks. + async fn replacement_pool_metadata_applies(&self, _opts: &HealOpts) -> Result { + Ok(true) + } + /// Read target-specific physical evidence for one replacement version. /// /// This is only used by automatic replacement healing after the normal @@ -1268,6 +1276,18 @@ impl HealStorageAPI for ECStoreHealStorage { .map_err(Error::Storage) } + async fn replacement_pool_metadata_applies(&self, opts: &HealOpts) -> Result { + let pool_index = opts + .pool + .ok_or_else(|| Error::other("replacement pool metadata is missing pool scope"))?; + let set_index = opts + .set + .ok_or_else(|| Error::other("replacement pool metadata is missing set scope"))?; + self.ecstore + .replacement_pool_metadata_applies(pool_index, set_index) + .map_err(Error::Storage) + } + async fn replacement_targets_have_version( &self, bucket: &str,