From 201e314dd9820712251b9f3b10c56ddb8fc06681 Mon Sep 17 00:00:00 2001 From: cxymds Date: Mon, 14 Sep 2026 21:04:19 +0800 Subject: [PATCH] fix(heal): scope nested walks and reject admin read-repair (#7851) * fix(heal): scope nested walks and reject admin read-repair * fix(ci): allocate stack for heal metadata regression --- .config/nextest.toml | 12 ++ crates/ecstore/src/set_disk/ops/heal_walk.rs | 187 ++++++++++++++++--- crates/protos/src/heal_control.rs | 9 + crates/protos/src/lib.rs | 14 +- rustfs/src/admin/handlers/heal.rs | 45 +++++ 5 files changed, 239 insertions(+), 28 deletions(-) diff --git a/.config/nextest.toml b/.config/nextest.toml index a600c4e3a..d26f7bc79 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -76,6 +76,12 @@ setup = 'ecstore-base-stack' filter = 'binary(lifecycle_integration_test) | (package(rustfs) & test(/^app::lifecycle_transition_api_test::/))' setup = 'lifecycle-large-stack' +# The pool metadata placement regression composes the full erasure-heal +# startup future and overflows the default Linux test-thread stack. +[[profile.default.scripts]] +filter = 'package(rustfs-heal) & test(=replacement_pool_metadata_follows_real_two_set_placement)' +setup = 'ecstore-large-stack' + [[profile.default.overrides]] filter = 'package(rustfs-ecstore) & (test(concurrent_resend_same_part_commits_one_generation) | test(concurrent_config_writes_from_separate_nodes_do_not_lose_writes) | test(/^store::bucket::tests::bucket_delete_(mark_delete|purge_removes|default_s3_delete)/))' test-group = 'ecstore-serial-flaky' @@ -239,6 +245,12 @@ setup = 'ecstore-base-stack' filter = 'binary(lifecycle_integration_test) | (package(rustfs) & test(/^app::lifecycle_transition_api_test::/))' setup = 'lifecycle-large-stack' +# The pool metadata placement regression composes the full erasure-heal +# startup future and overflows the default Linux test-thread stack. +[[profile.ci.scripts]] +filter = 'package(rustfs-heal) & test(=replacement_pool_metadata_follows_real_two_set_placement)' +setup = 'ecstore-large-stack' + # =========================================================================== # QUARANTINE — flaky tests granted retries = 2 under the ci profile ONLY. # diff --git a/crates/ecstore/src/set_disk/ops/heal_walk.rs b/crates/ecstore/src/set_disk/ops/heal_walk.rs index 2dbf2fdd5..aaa529850 100644 --- a/crates/ecstore/src/set_disk/ops/heal_walk.rs +++ b/crates/ecstore/src/set_disk/ops/heal_walk.rs @@ -72,6 +72,12 @@ struct HealWalkObject { /// order, so each successful ingest corresponds to exactly one new object. struct HealWalkCollector { bucket: String, + /// Full S3 prefix used after moving the walk to a non-root base directory. + /// The local walker emits the base directory marker before applying its + /// relative filter, so the boundary must also be checked on the decoded key. + prefix: String, + /// Inclusive full object-key cursor for the current page. + forward_to: Option, batch_objects: usize, version_budget: usize, include_lifecycle_object_info: bool, @@ -83,6 +89,10 @@ struct HealWalkCollector { } impl HealWalkCollector { + fn accepts_name(&self, name: &str) -> bool { + name.starts_with(&self.prefix) && self.forward_to.as_deref().is_none_or(|forward_to| name >= forward_to) + } + fn lock_objects(&self) -> disk::error::Result>> { self.objects.lock().map_err(|_| { self.cancel.cancel(); @@ -99,6 +109,29 @@ impl HealWalkCollector { self.cancel.cancel(); } + fn record_object(&self, name: String, versions: Vec) -> Option<(usize, usize)> { + let mut objects = self.lock_objects().ok()?; + let mut added = 0; + if let Some(existing) = objects.iter_mut().find(|object| object.name == name) { + let mut seen: HashSet = existing + .versions + .iter() + .filter_map(|version| version.version_id.clone()) + .collect(); + for version in versions { + if seen.insert(version.version_id.clone().unwrap_or_default()) { + existing.versions.push(version); + added += 1; + } + } + } else { + added = versions.len(); + objects.push(HealWalkObject { name, versions }); + } + let ver_total = self.version_total.fetch_add(added, Ordering::SeqCst) + added; + Some((objects.len(), ver_total)) + } + fn take_decode_error(&self) -> disk::error::Result> { self.decode_error.lock().map(|mut error| error.take()).map_err(|_| { self.cancel.cancel(); @@ -111,6 +144,10 @@ impl HealWalkCollector { /// is met — always at a sorted object-key boundary so a heavily-versioned /// object is never split across pages. fn ingest(&self, entry: MetaCacheEntry) { + if !self.accepts_name(&entry.name) { + return; + } + // Skip pure directory entries; they carry no versions to heal here. if entry.is_dir() { return; @@ -147,18 +184,8 @@ impl HealWalkCollector { return; } - let added = versions.len(); - let (objs_len, ver_total) = { - let Ok(mut objects) = self.lock_objects() else { - return; - }; - objects.push(HealWalkObject { - name: entry.name, - versions, - }); - let objs_len = objects.len(); - let ver_total = self.version_total.fetch_add(added, Ordering::SeqCst) + added; - (objs_len, ver_total) + let Some((objs_len, ver_total)) = self.record_object(entry.name, versions) else { + return; }; // Bound at an object boundary (this object is fully included). @@ -178,7 +205,7 @@ impl HealWalkCollector { let mut versions = Vec::new(); for entry in entries.0.iter().flatten() { - if entry.is_dir() || entry.name.is_empty() { + if !self.accepts_name(&entry.name) || entry.is_dir() || entry.name.is_empty() { continue; } if name.is_empty() { @@ -215,15 +242,8 @@ impl HealWalkCollector { return; } - let added = versions.len(); - let (objs_len, ver_total) = { - let Ok(mut objects) = self.lock_objects() else { - return; - }; - objects.push(HealWalkObject { name, versions }); - let objs_len = objects.len(); - let ver_total = self.version_total.fetch_add(added, Ordering::SeqCst) + added; - (objs_len, ver_total) + let Some((objs_len, ver_total)) = self.record_object(name, versions) else { + return; }; if objs_len >= self.batch_objects || ver_total >= self.version_budget { @@ -290,6 +310,8 @@ impl SetDisks { let collector = Arc::new(HealWalkCollector { bucket: bucket.to_string(), + prefix: prefix.to_string(), + forward_to: forward_to.map(str::to_owned), batch_objects, version_budget: version_budget.max(1), include_lifecycle_object_info, @@ -303,12 +325,26 @@ impl SetDisks { let agreed_collector = collector.clone(); let partial_collector = collector.clone(); - let filter_prefix = if prefix.is_empty() { None } else { Some(prefix.to_string()) }; + // `scan_dir` applies `filter_prefix` to names relative to `path`. Keep + // the full prefix in the collector as a second boundary check because + // the walker emits the base directory marker before that filter. + let path = rustfs_utils::path::base_dir_from_prefix(prefix); + let filter_prefix = if prefix.is_empty() { + None + } else { + Some( + prefix + .trim_start_matches(&path) + .trim_start_matches('/') + .trim_end_matches('/') + .to_owned(), + ) + }; let opts = ListPathRawOptions { disks, bucket: bucket.to_string(), - path: String::new(), + path, recursive: true, incl_deleted: true, filter_prefix, @@ -372,8 +408,14 @@ mod tests { use uuid::Uuid; fn test_collector() -> Arc { + test_collector_for("", None) + } + + fn test_collector_for(prefix: &str, forward_to: Option<&str>) -> Arc { Arc::new(HealWalkCollector { bucket: "bucket".to_string(), + prefix: prefix.to_string(), + forward_to: forward_to.map(str::to_owned), batch_objects: 2, version_budget: 2, include_lifecycle_object_info: false, @@ -385,6 +427,101 @@ mod tests { }) } + #[test] + fn collector_rechecks_full_prefix_and_forward_cursor() { + let collector = test_collector_for("a/b/", Some("a/b/child")); + + assert!(!collector.accepts_name("a/"), "base marker outside the requested prefix must be ignored"); + assert!(!collector.accepts_name("a/b-other/object"), "sibling prefixes must not be included"); + assert!( + !collector.accepts_name("a/b/before"), + "objects before the inclusive cursor must be ignored" + ); + assert!(collector.accepts_name("a/b/child"), "the cursor boundary is inclusive"); + assert!(collector.accepts_name("a/b/child/deep"), "descendants after the cursor must be included"); + } + + #[tokio::test] + async fn heal_walk_scopes_non_root_prefix_to_nested_objects() { + use rustfs_filemeta::FileInfo; + + let (temp_dirs, disks, set_disks) = hermetic_set_disks_isolated(4).await; + let bucket = "nested-prefix-bucket"; + for disk in &disks { + disk.make_volume(bucket).await.expect("test bucket should be created"); + } + let objects = [ + "a/", + "a/b", + "a/b/direct", + "a/b/child/deep", + "a/b/child/later", + "a/b-other/sibling", + "a/c/outside", + "single", + "中文/%2F+ key/item", + ]; + + // A single replica is below read quorum. The raw UNION must still + // discover every version, including explicit directory markers. + for (index, object) in objects.iter().enumerate() { + let mut info = FileInfo::new(object, 4, 2); + info.volume = bucket.to_string(); + info.name = object.to_string(); + info.version_id = Some(Uuid::from_u128(u128::try_from(index + 1).expect("index fits UUID"))); + info.versioned = true; + info.size = if object.ends_with('/') { 0 } else { 100 }; + info.mod_time = Some( + OffsetDateTime::from_unix_timestamp(i64::try_from(index + 1).expect("index fits timestamp")) + .expect("fixture timestamp"), + ); + let mut metadata = FileMeta::new(); + metadata + .add_version(info) + .expect("fixture metadata should accept the version"); + let object_dir = temp_dirs[0] + .path() + .join(bucket) + .join(rustfs_utils::path::encode_dir_object(object)); + tokio::fs::create_dir_all(&object_dir) + .await + .expect("object directory should be created"); + tokio::fs::write( + object_dir.join(crate::disk::STORAGE_FORMAT_FILE), + metadata.marshal_msg().expect("fixture metadata should serialize"), + ) + .await + .expect("fixture metadata should be written"); + } + + for prefix in ["", "a/", "a/b", "a/b/", "single", "中文/", "中文/%2F+", "absent/"] { + let mut expected: Vec<_> = objects.iter().copied().filter(|object| object.starts_with(prefix)).collect(); + expected.sort_unstable(); + let mut actual = Vec::new(); + let mut forward: Option = None; + let mut pages = 0; + loop { + let (versions, next, truncated) = set_disks + .heal_walk_versions_page(bucket, prefix, forward.as_deref(), 2, 100, false) + .await + .expect("prefix disk walk should succeed"); + pages += 1; + actual.extend(versions.into_iter().map(|version| version.name)); + if !truncated { + assert!(next.is_none(), "complete prefix page must clear its cursor"); + break; + } + let next = next.expect("truncated prefix page must return a cursor"); + assert!(next.starts_with(prefix), "cursor must retain the full logical key"); + assert!(forward.as_ref().is_none_or(|previous| &next > previous), "cursor must advance"); + forward = Some(next); + assert!(pages < 20, "prefix pagination must terminate"); + } + actual.sort_unstable(); + assert_eq!(actual, expected, "prefix {prefix:?} must enumerate every matching key exactly once"); + } + } + #[test] fn collectors_preserve_historical_null_identity() { use rustfs_filemeta::FileInfo; @@ -573,6 +710,8 @@ mod tests { let collector = Arc::new(HealWalkCollector { bucket: "bucket".to_string(), + prefix: "".to_string(), + forward_to: None, batch_objects: 1000, version_budget: 10_000, include_lifecycle_object_info: false, diff --git a/crates/protos/src/heal_control.rs b/crates/protos/src/heal_control.rs index b3d7c8282..d197349a3 100644 --- a/crates/protos/src/heal_control.rs +++ b/crates/protos/src/heal_control.rs @@ -861,6 +861,15 @@ mod tests { let unknown = rmp_serde::to_vec_named(&unknown).unwrap(); assert!(decode_envelope(&unknown).unwrap_err().contains("unknown field")); + let mut read_repair = serde_json::to_value(&envelope).expect("start envelope should serialize"); + read_repair["command"]["request"] + .as_object_mut() + .expect("start request must be an object") + .insert("readRepair".to_string(), serde_json::Value::Bool(true)); + let encoded = rmp_serde::to_vec_named(&read_repair).expect("invalid start fixture should encode"); + let error = decode_envelope(&encoded).expect_err("RPC must reject an Admin readRepair field"); + assert!(error.contains("unknown field") && error.contains("readRepair"), "{error}"); + let executable = Envelope::start(test_request(request_id.clone()), RequestMetadata::new([1; 16], 10_000, 20_000, 7)).unwrap(); assert!(executable.validate_execution(15_000, 7).is_ok()); diff --git a/crates/protos/src/lib.rs b/crates/protos/src/lib.rs index 82d301124..fc806c006 100644 --- a/crates/protos/src/lib.rs +++ b/crates/protos/src/lib.rs @@ -172,7 +172,11 @@ pub const HEAL_CONTROL_RPC_MAX_MESSAGE_SIZE: usize = heal_control::RESULT_MAX_SI pub const HEAL_CONTROL_PROTOCOL_VERSION: u32 = 3; pub const DYNAMIC_CONFIG_PROTOCOL_VERSION: u32 = 1; pub const BACKGROUND_HEAL_STATUS_PROTOCOL_VERSION: u32 = 2; -pub const HEAL_CONTROL_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-heal-control-capability-v3\0"; +// v4 is an admission boundary, not a transport version bump. A peer that +// cannot recognize this probe must fail the Admin Heal capability preflight; +// accepting the older v3 probe would allow an upgraded node to silently lose +// newly validated request semantics during a rolling upgrade. +pub const HEAL_CONTROL_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-heal-control-capability-v4\0"; pub const REMOTE_VERSION_STATE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-tier-remote-version-state-capability-v1\0"; pub const CROSS_POOL_FENCE_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-cross-pool-fence-capability-v1\0"; pub const ILM_RECOVERY_EXPORT_CAPABILITY_PROBE_PREFIX: &[u8] = b"rustfs-ilm-recovery-export-capability-v1\0"; @@ -318,7 +322,7 @@ pub fn canonical_heal_control_capability_ack( topology_fingerprint: &str, probe: &[u8], ) -> Result, std::num::TryFromIntError> { - const DOMAIN: &[u8] = b"rustfs-heal-control-capability-ack-v3\0"; + const DOMAIN: &[u8] = b"rustfs-heal-control-capability-ack-v4\0"; let fingerprint = topology_fingerprint.as_bytes(); let mut body = Vec::with_capacity(DOMAIN.len() + 4 + 8 + fingerprint.len() + 8 + probe.len()); @@ -2261,10 +2265,10 @@ mod heal_control_tests { #[test] fn canonical_capability_ack_binds_version_and_topology() { assert_eq!(HEAL_CONTROL_PROTOCOL_VERSION, 3); - assert!(HEAL_CONTROL_CAPABILITY_PROBE_PREFIX.starts_with(b"rustfs-heal-control-capability-v3")); + assert!(HEAL_CONTROL_CAPABILITY_PROBE_PREFIX.starts_with(b"rustfs-heal-control-capability-v4")); let probe = heal_control_capability_probe(&[7; 16]); let ack = canonical_heal_control_capability_ack(1, "ab", &probe).expect("small acknowledgement should encode"); - let mut golden = b"rustfs-heal-control-capability-ack-v3\0".to_vec(); + let mut golden = b"rustfs-heal-control-capability-ack-v4\0".to_vec(); golden.extend_from_slice(&1_u32.to_be_bytes()); golden.extend_from_slice(&2_u64.to_be_bytes()); golden.extend_from_slice(b"ab"); @@ -2279,6 +2283,8 @@ mod heal_control_tests { ); assert!(is_heal_control_capability_probe(&probe)); assert!(!is_heal_control_capability_probe(HEAL_CONTROL_CAPABILITY_PROBE_PREFIX)); + let legacy_probe = [b"rustfs-heal-control-capability-v3\0".as_slice(), &[7; 16]].concat(); + assert!(!is_heal_control_capability_probe(&legacy_probe)); } #[test] diff --git a/rustfs/src/admin/handlers/heal.rs b/rustfs/src/admin/handlers/heal.rs index 8148e4c01..021a27291 100644 --- a/rustfs/src/admin/handlers/heal.rs +++ b/rustfs/src/admin/handlers/heal.rs @@ -1039,6 +1039,7 @@ where E: FnOnce() -> EF, EF: Future>, { + validate_heal_start_options(options)?; validate_heal_selector(endpoints, options.pool, options.set) .map_err(|err| admin_error(S3ErrorCode::InvalidArgument, err.to_string()))?; probe().await?; @@ -1346,6 +1347,17 @@ fn validate_heal_request_mode(hip: &HealInitParams) -> S3Result<()> { Ok(()) } +fn validate_heal_start_options(options: &HealOpts) -> S3Result<()> { + if options.read_repair { + return Err(admin_error( + S3ErrorCode::InvalidArgument, + "readRepair=true is not supported for Admin Heal", + )); + } + + Ok(()) +} + fn json_response(status: StatusCode, body: Vec) -> S3Response<(StatusCode, Body)> { let mut headers = HeaderMap::new(); headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json")); @@ -1451,6 +1463,9 @@ impl Operation for HealHandler { }; let hip = extract_heal_init_params(&bytes, &req.uri, params)?; validate_heal_request_mode(&hip)?; + if hip.client_token.is_empty() && !hip.force_stop { + validate_heal_start_options(&hip.hs)?; + } let response_operation = if hip.force_stop { "cancel_heal" } else if !hip.client_token.is_empty() && !hip.force_start { @@ -1906,6 +1921,36 @@ mod tests { assert!(executed.load(Ordering::SeqCst)); } + #[tokio::test] + async fn read_repair_admin_start_is_rejected_before_probe_or_execution() { + let probed = AtomicBool::new(false); + let executed = AtomicBool::new(false); + let options = HealOpts { + read_repair: true, + ..Default::default() + }; + + let error = execute_after_heal_start_preflight( + &super::EndpointServerPools::default(), + &options, + || async { + probed.store(true, Ordering::SeqCst); + Ok(()) + }, + || async { + executed.store(true, Ordering::SeqCst); + Ok(()) + }, + ) + .await + .expect_err("Admin Heal must reject the internal read-repair mode"); + + assert_eq!(error.code(), &S3ErrorCode::InvalidArgument); + assert_eq!(error.code().status_code(), Some(StatusCode::BAD_REQUEST)); + assert!(!probed.load(Ordering::SeqCst), "rejected options must precede capability probing"); + assert!(!executed.load(Ordering::SeqCst), "rejected options must not execute a heal"); + } + #[tokio::test] async fn heal_start_retry_preflight_failures_do_not_create_request_identities() { let hip = HealInitParams {