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
This commit is contained in:
cxymds
2026-09-14 21:04:19 +08:00
committed by GitHub
parent 28db20f7c4
commit 201e314dd9
5 changed files with 239 additions and 28 deletions
+12
View File
@@ -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.
#
+163 -24
View File
@@ -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<String>,
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<std::sync::MutexGuard<'_, Vec<HealWalkObject>>> {
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<HealWalkVersion>) -> 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<String> = 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<Option<DiskError>> {
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<HealWalkCollector> {
test_collector_for("", None)
}
fn test_collector_for(prefix: &str, forward_to: Option<&str>) -> Arc<HealWalkCollector> {
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<String> = 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,
+9
View File
@@ -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());
+10 -4
View File
@@ -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<Vec<u8>, 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]
+45
View File
@@ -1039,6 +1039,7 @@ where
E: FnOnce() -> EF,
EF: Future<Output = S3Result<T>>,
{
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<u8>) -> 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 {