Files
rustfs/crates/ecstore/src/set_disk/metadata.rs
T
Zhengchao An a6a04b5faa refactor(ecstore,rustfs): reuse canonical starts_with_ignore_ascii_case (#6759)
* refactor(ecstore,rustfs): reuse canonical starts_with_ignore_ascii_case

`crates/utils/src/http/metadata_compat.rs` owns the internal metadata key helpers, including `starts_with_ignore_ascii_case`. Two files carried their own byte-identical copies of that predicate: `SetDisks::starts_with_ignore_ascii_case` in ecstore and a free function in the S3 options layer. Both drive internal metadata key classification (`internal_metadata_suffix` and quorum hashing on one side, `should_skip_object_metadata_key` and `is_reserved_user_metadata_key` on the other), so keeping three implementations of one predicate is an avoidable drift risk on a path that decides whether an internal key is treated as user metadata.

Delete both local copies and call the canonical implementation. Every prefix used at these call sites is an ASCII constant or literal, where the canonical byte-slice comparison and the removed `str::get(..n)` form are equivalent; that equivalence was checked differentially over 4.6M (key, prefix) pairs, including keys with multi-byte characters straddling the prefix boundary. No other logic in `internal_metadata_suffix` or `should_skip_object_metadata_key` changed.

Add regression tests on both sides pinning the two properties the switch depends on: internal prefixes match case-insensitively (a mixed-case `X-RustFS-Internal-*` key stays internal), and keys shorter than a prefix never match (they stay ordinary user metadata).

Refs rustfs/backlog#2051

* fix(rustfs): avoid typos-checker false positive in prefix-length test

The test literal "x-rustfs-encryptio" (a deliberate truncation of the
x-rustfs-encryption- prefix, used to assert that a key shorter than every
internal prefix falls through to user metadata) reads as a likely typo of
"encryption" to the repo's typos CI check. Derive it from
RUSTFS_ENCRYPTION_PREFIX via slicing instead of a hand-typed literal, which
both satisfies the linter and ties the truncation to the real constant
instead of a copy-typed guess.

Refs rustfs/backlog#2051
2026-08-28 07:35:09 +08:00

1673 lines
65 KiB
Rust

// Copyright 2024 RustFS Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use super::{
Bytes, DATA_MOVEMENT_MULTIPART_PREFIX, DiskError, DiskStore, FileInfo, HashMap, HashSet, OBJECT_OP_IGNORED_ERRS, ObjProps,
OffsetDateTime, SetDisks, Sha256, TRANSITION_COMPLETE, Uuid, debug, disk, error, file_info_is_valid_for_metadata, hex,
reduce_read_quorum_errs, warn,
};
#[cfg(test)]
use crate::disk::DiskOption;
#[cfg(test)]
use crate::disk::endpoint::Endpoint;
#[cfg(test)]
use crate::disk::new_disk;
use rustfs_utils::http;
use sha2::Digest;
#[derive(Clone, Copy)]
struct FileInfoIdentityGroup {
hash: [u8; 32],
count: usize,
mod_time: Option<OffsetDateTime>,
}
impl SetDisks {
pub(super) fn all_not_found_metadata(errs: &[Option<DiskError>]) -> bool {
!errs.is_empty()
&& errs.iter().all(|err| match err {
Some(err) => {
matches!(
err,
DiskError::FileNotFound
| DiskError::FileVersionNotFound
| DiskError::VolumeNotFound
| DiskError::DiskNotFound
) || OBJECT_OP_IGNORED_ERRS.contains(err)
}
None => false,
})
&& errs.iter().any(|err| {
matches!(
err,
Some(DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::VolumeNotFound)
)
})
}
pub(super) fn reduce_common_data_dir(data_dirs: &[Option<Uuid>], write_quorum: usize) -> Option<Uuid> {
let mut data_dirs_count = HashMap::new();
for ddir in data_dirs.iter().flatten().copied() {
*data_dirs_count.entry(ddir).or_insert(0) += 1;
}
let mut max = 0;
let mut data_dir = None;
for (ddir, count) in data_dirs_count {
if count > max {
max = count;
data_dir = Some(ddir);
}
}
if max >= write_quorum { data_dir } else { None }
}
pub(super) fn get_upload_id_dir(bucket: &str, object: &str, upload_id: &str) -> String {
let upload_uuid = base64_simd::URL_SAFE_NO_PAD
.decode_to_vec(upload_id.as_bytes())
.and_then(|v| {
String::from_utf8(v).map_or_else(
|_| Ok(upload_id.to_owned()),
|v| {
let parts: Vec<_> = v.splitn(2, '.').collect();
if parts.len() == 2 {
Ok(parts[1].to_string())
} else {
Ok(upload_id.to_string())
}
},
)
})
.unwrap_or_default();
format!("{}/{}", Self::get_multipart_sha_dir(bucket, object), upload_uuid)
}
pub(super) fn get_multipart_upload_dir(bucket: &str, object: &str, upload_id: &str, data_movement: bool) -> String {
let upload_dir = Self::get_upload_id_dir(bucket, object, upload_id);
if data_movement {
format!("{DATA_MOVEMENT_MULTIPART_PREFIX}/{upload_dir}")
} else {
upload_dir
}
}
pub(super) fn get_multipart_sha_dir(bucket: &str, object: &str) -> String {
let path = format!("{bucket}/{object}");
let mut hasher = Sha256::new();
hasher.update(path);
hex(hasher.finalize())
}
pub(super) fn common_parity(parities: &[i32], default_parity_count: i32) -> i32 {
let n = parities.len() as i32;
let mut occ_map: HashMap<i32, i32> = HashMap::new();
for &p in parities {
*occ_map.entry(p).or_insert(0) += 1;
}
let mut max_occ = 0;
let mut cparity = 0;
for (&parity, &occ) in &occ_map {
if parity == -1 {
// Ignore non defined parity
continue;
}
let mut read_quorum = n - parity;
if default_parity_count > 0 && parity == 0 {
// In this case, parity == 0 implies that this object version is a
// delete marker
read_quorum = n / 2 + 1;
}
if occ < read_quorum {
// Ignore this parity since we don't have enough shards for read quorum
continue;
}
if occ > max_occ {
max_occ = occ;
cparity = parity;
}
}
if max_occ == 0 {
// Did not find anything useful
return -1;
}
cparity
}
pub(super) fn list_object_modtimes(parts_metadata: &[FileInfo], errs: &[Option<DiskError>]) -> Vec<Option<OffsetDateTime>> {
let mut times = vec![None; parts_metadata.len()];
for (i, metadata) in parts_metadata.iter().enumerate() {
if errs[i].is_some() {
continue;
}
times[i] = metadata.mod_time
}
times
}
pub(super) fn common_time(times: &[Option<OffsetDateTime>], quorum: usize) -> Option<OffsetDateTime> {
let (time, count) = Self::common_time_and_occurrence(times);
if count >= quorum { time } else { None }
}
pub(super) fn common_time_and_occurrence(times: &[Option<OffsetDateTime>]) -> (Option<OffsetDateTime>, usize) {
let mut time_occurrence_map = HashMap::new();
// Ignore the uuid sentinel and count the rest.
for time in times.iter().flatten() {
*time_occurrence_map.entry(time.unix_timestamp_nanos()).or_insert(0) += 1;
}
let mut maxima = 0; // Counter for remembering max occurrence of elements.
let mut latest = 0;
// Find the common cardinality from previously collected
// occurrences of elements.
for (&nano, &count) in &time_occurrence_map {
if count < maxima {
continue;
}
// We are at or above maxima
if count > maxima || nano > latest {
maxima = count;
latest = nano;
}
}
if latest == 0 {
return (None, maxima);
}
if let Ok(time) = OffsetDateTime::from_unix_timestamp_nanos(latest) {
(Some(time), maxima)
} else {
(None, maxima)
}
}
pub(super) fn common_etag(etags: &[Option<String>], quorum: usize) -> Option<String> {
let (etag, count) = Self::common_etags(etags);
if count >= quorum { etag } else { None }
}
pub(super) fn common_etags(etags: &[Option<String>]) -> (Option<String>, usize) {
let mut etags_map = HashMap::new();
for etag in etags.iter().flatten() {
*etags_map.entry(etag).or_insert(0) += 1;
}
let mut maxima = 0; // Counter for remembering max occurrence of elements.
let mut latest = None;
for (&etag, &count) in &etags_map {
if count < maxima {
continue;
}
// We are at or above maxima
if count > maxima {
maxima = count;
latest = Some(etag.clone());
}
}
(latest, maxima)
}
pub(super) fn list_object_etags(parts_metadata: &[FileInfo], errs: &[Option<DiskError>]) -> Vec<Option<String>> {
let mut etags = vec![None; parts_metadata.len()];
for (i, metadata) in parts_metadata.iter().enumerate() {
if errs[i].is_some() {
continue;
}
if let Some(etag) = metadata.metadata.get("etag") {
etags[i] = Some(etag.clone())
}
}
etags
}
pub(super) fn list_object_parities(parts_metadata: &[FileInfo], errs: &[Option<DiskError>]) -> Vec<i32> {
let total_shards = parts_metadata.len();
let total_shards_i32 = i32::try_from(total_shards).unwrap_or(i32::MAX);
let half = total_shards_i32 / 2;
let mut parities: Vec<i32> = vec![-1; total_shards];
for (index, metadata) in parts_metadata.iter().enumerate() {
if errs[index].is_some() {
parities[index] = -1;
continue;
}
if !file_info_is_valid_for_metadata(metadata) {
parities[index] = -1;
continue;
}
// A parity count outside [0, total_shards] cannot describe a real
// layout on this set: it comes from corrupt or foreign metadata
// (e.g. stray leftovers, rustfs#5801). Treat the entry as invalid
// instead of clamping to i32::MAX, which would poison
// `common_parity`'s occurrence counting.
let erasure_parity = i32::try_from(metadata.erasure.parity_blocks).unwrap_or(-1);
let erasure_parity = if (0..=total_shards_i32).contains(&erasure_parity) {
erasure_parity
} else {
-1
};
if metadata.is_canonical_delete_marker() || metadata.size == 0 {
parities[index] = half;
} else if erasure_parity < 0 {
parities[index] = -1;
} else if metadata.transition_status == TRANSITION_COMPLETE {
let majority_metadata_parity = total_shards_i32 - (half + 1);
parities[index] = majority_metadata_parity.max(erasure_parity);
} else {
parities[index] = erasure_parity;
}
}
parities
}
#[tracing::instrument(level = "debug", skip(parts_metadata))]
pub(super) fn object_quorum_from_meta(
parts_metadata: &[FileInfo],
errs: &[Option<DiskError>],
default_parity_count: usize,
) -> disk::error::Result<(i32, i32)> {
if Self::all_not_found_metadata(errs) {
return Err(DiskError::FileNotFound);
}
let expected_rquorum = if default_parity_count == 0 {
parts_metadata.len()
} else {
parts_metadata.len() / 2
};
if let Some(err) = reduce_read_quorum_errs(errs, OBJECT_OP_IGNORED_ERRS, expected_rquorum) {
// let object = parts_metadata.first().map(|v| v.name.clone()).unwrap_or_default();
// error!("object_quorum_from_meta: {:?}, errs={:?}, object={:?}", err, errs, object);
return Err(err);
}
if default_parity_count == 0 {
return Ok((parts_metadata.len() as i32, parts_metadata.len() as i32));
}
let parities = Self::list_object_parities(parts_metadata, errs);
let parity_blocks = Self::common_parity(&parities, default_parity_count as i32);
if parity_blocks < 0 {
// No parity value reached read quorum. Distinguish two cases:
// enough disks answered with valid-looking metadata that simply
// cannot be reconciled (corrupt/foreign entries — retrying cannot
// help, and heal should see Corrupt, rustfs#5801) versus too few
// healthy answers (a genuine quorum condition where retry may
// succeed once disks recover).
let healthy_replies = errs.iter().filter(|err| err.is_none()).count();
if healthy_replies >= expected_rquorum {
error!(
"object_quorum_from_meta: irreconcilable parity across {healthy_replies} healthy replies (corrupt metadata), errs={errs:?}"
);
return Err(DiskError::FileCorrupt);
}
error!("object_quorum_from_meta: parity_blocks < 0, errs={:?}", errs);
return Err(DiskError::ErasureReadQuorum);
}
let data_blocks = parts_metadata.len() as i32 - parity_blocks;
let write_quorum = if data_blocks == parity_blocks {
data_blocks + 1
} else {
data_blocks
};
Ok((data_blocks, write_quorum))
}
#[tracing::instrument(level = "debug", skip(disks, parts_metadata))]
pub(super) fn list_online_disks(
disks: &[Option<DiskStore>],
parts_metadata: &[FileInfo],
errs: &[Option<DiskError>],
quorum: usize,
) -> (Vec<Option<DiskStore>>, Option<OffsetDateTime>, Option<String>) {
let mod_times = Self::list_object_modtimes(parts_metadata, errs);
let mod_time = Self::common_time(&mod_times, quorum);
if mod_time.is_none() {
let etags = Self::list_object_etags(parts_metadata, errs);
let etag_op = Self::common_etag(&etags, quorum);
if let Some(etag) = etag_op {
let mut new_disk = vec![None; disks.len()];
for (i, etag_item) in etags.iter().enumerate() {
if let Some(etag_item) = etag_item
&& etag_item == &etag
&& file_info_is_valid_for_metadata(&parts_metadata[i])
{
new_disk[i].clone_from(&disks[i]);
}
}
return (new_disk, None, Some(etag));
}
}
let mut new_disk = vec![None; disks.len()];
for (i, &t) in mod_times.iter().enumerate() {
if file_info_is_valid_for_metadata(&parts_metadata[i]) && mod_time == t {
new_disk[i].clone_from(&disks[i]);
}
}
(new_disk, mod_time, None)
}
fn usable_fileinfo_count(parts_metadata: &[FileInfo], errs: &[Option<DiskError>]) -> (usize, bool) {
let mut has_read_error = false;
let mut usable_metadata = 0;
for (meta, err) in parts_metadata.iter().zip(errs.iter()) {
if err.is_some() {
has_read_error = true;
continue;
}
if file_info_is_valid_for_metadata(meta) {
usable_metadata += 1;
}
}
(usable_metadata, has_read_error)
}
pub(super) fn latest_fileinfo_selection_quorum(
version_id: &str,
parts_metadata: &[FileInfo],
errs: &[Option<DiskError>],
read_quorum: usize,
write_quorum: usize,
) -> usize {
if !version_id.is_empty() || write_quorum <= read_quorum {
return read_quorum;
}
let (usable_metadata, has_read_error) = Self::usable_fileinfo_count(parts_metadata, errs);
if usable_metadata < write_quorum {
return read_quorum;
}
if !has_read_error {
return write_quorum;
}
let mut identity_counts = HashMap::with_capacity(usable_metadata);
for (meta, err) in parts_metadata.iter().zip(errs.iter()) {
if err.is_some() || !file_info_is_valid_for_metadata(meta) {
continue;
}
let key = Self::file_info_quorum_hash(meta);
let count = identity_counts.entry(key).or_insert(0);
*count += 1;
if *count >= write_quorum {
return write_quorum;
}
}
read_quorum
}
pub(super) fn select_valid_fileinfo(
disks: &[Option<DiskStore>],
parts_metadata: &[FileInfo],
errs: &[Option<DiskError>],
version_id: &str,
read_quorum: usize,
write_quorum: usize,
) -> disk::error::Result<(Vec<Option<DiskStore>>, FileInfo, usize)> {
let selection_quorum =
Self::latest_fileinfo_selection_quorum(version_id, parts_metadata, errs, read_quorum, write_quorum);
let (usable_metadata, has_read_error) = Self::usable_fileinfo_count(parts_metadata, errs);
if version_id.is_empty()
&& write_quorum > read_quorum
&& has_read_error
&& usable_metadata >= write_quorum
&& selection_quorum == read_quorum
{
let (online_disks, fi) = Self::pick_degraded_latest_fileinfo(disks, parts_metadata, errs, read_quorum, write_quorum)?;
return Ok((online_disks, fi, read_quorum));
}
let (online_disks, mod_time, etag) = Self::list_online_disks(disks, parts_metadata, errs, selection_quorum);
let fi = Self::pick_valid_fileinfo(parts_metadata, mod_time, etag, selection_quorum)?;
Ok((online_disks, fi, selection_quorum))
}
pub(super) fn pick_valid_fileinfo(
metas: &[FileInfo],
mod_time: Option<OffsetDateTime>,
etag: Option<String>,
quorum: usize,
) -> disk::error::Result<FileInfo> {
Self::find_file_info_in_quorum(metas, &mod_time, &etag, quorum)
}
pub(crate) fn hydrate_selected_fileinfo_part_checksums(fi: &mut FileInfo) -> disk::error::Result<()> {
fi.hydrate_data_movement_part_checksums().map_err(DiskError::from)?;
for part in &fi.parts {
let Some(checksums) = part.checksums.as_ref() else {
continue;
};
let mut algorithms = HashSet::with_capacity(checksums.len());
for (name, value) in checksums {
let Some(checksum) = rustfs_rio::Checksum::new_from_string(name, value) else {
return Err(DiskError::FileCorrupt);
};
if checksum.checksum_type.is(rustfs_rio::ChecksumType::MULTIPART) {
return Err(DiskError::FileCorrupt);
}
if !algorithms.insert(checksum.checksum_type.base().0) {
return Err(DiskError::FileCorrupt);
}
}
}
Ok(())
}
fn update_hash_bytes(hasher: &mut Sha256, value: &[u8]) {
hasher.update(value.len().to_le_bytes());
hasher.update(value);
}
fn update_hash_str(hasher: &mut Sha256, value: &str) {
Self::update_hash_bytes(hasher, value.as_bytes());
}
fn update_hash_optional_uuid(hasher: &mut Sha256, value: Option<Uuid>) {
if let Some(value) = value {
hasher.update([1]);
hasher.update(value.as_bytes());
} else {
hasher.update([0]);
}
}
fn update_hash_optional_time(hasher: &mut Sha256, value: Option<OffsetDateTime>) {
if let Some(value) = value {
hasher.update([1]);
hasher.update(value.unix_timestamp_nanos().to_le_bytes());
} else {
hasher.update([0]);
}
}
fn update_hash_optional_u32(hasher: &mut Sha256, value: Option<u32>) {
if let Some(value) = value {
hasher.update([1]);
hasher.update(value.to_le_bytes());
} else {
hasher.update([0]);
}
}
fn update_hash_optional_u64(hasher: &mut Sha256, value: Option<u64>) {
if let Some(value) = value {
hasher.update([1]);
hasher.update(value.to_le_bytes());
} else {
hasher.update([0]);
}
}
fn update_hash_optional_bytes(hasher: &mut Sha256, value: Option<&Bytes>) {
if let Some(value) = value {
hasher.update([1]);
Self::update_hash_bytes(hasher, value);
} else {
hasher.update([0]);
}
}
fn update_hash_optional_str(hasher: &mut Sha256, value: Option<&str>) {
if let Some(value) = value {
hasher.update([1]);
Self::update_hash_str(hasher, value);
} else {
hasher.update([0]);
}
}
fn file_info_has_encryption_metadata(meta: &FileInfo) -> bool {
meta.metadata.keys().any(|name| http::is_object_encryption_marker(name))
}
fn internal_metadata_suffix(name: &str) -> Option<&str> {
name.get(http::RUSTFS_INTERNAL_PREFIX.len()..)
.filter(|_| http::starts_with_ignore_ascii_case(name, http::RUSTFS_INTERNAL_PREFIX))
.or_else(|| {
name.get(http::MINIO_INTERNAL_PREFIX.len()..)
.filter(|_| http::starts_with_ignore_ascii_case(name, http::MINIO_INTERNAL_PREFIX))
})
}
fn is_replication_quorum_metadata_key(name: &str) -> bool {
if name.eq_ignore_ascii_case(http::AMZ_BUCKET_REPLICATION_STATUS) {
return true;
}
let Some(suffix) = Self::internal_metadata_suffix(name) else {
return false;
};
suffix.eq_ignore_ascii_case(http::SUFFIX_REPLICA_STATUS)
|| suffix.eq_ignore_ascii_case(http::SUFFIX_REPLICA_TIMESTAMP)
|| suffix.eq_ignore_ascii_case(http::SUFFIX_REPLICATION_STATUS)
|| suffix.eq_ignore_ascii_case(http::SUFFIX_REPLICATION_TIMESTAMP)
|| suffix.eq_ignore_ascii_case(http::SUFFIX_PURGESTATUS)
|| http::starts_with_ignore_ascii_case(suffix, http::SUFFIX_REPLICATION_RESET_ARN_PREFIX)
// Raw compatibility keys are normalized and hashed separately below.
|| http::starts_with_ignore_ascii_case(suffix, http::SUFFIX_REPLICATION_DELETE_MARKER_VERSION_ARN_PREFIX)
}
fn update_hash_quorum_metadata_map(hasher: &mut Sha256, entries: &HashMap<String, String>) {
let mut entries = entries
.iter()
.filter(|(name, _)| !Self::is_replication_quorum_metadata_key(name))
.collect::<Vec<_>>();
entries.sort_by(|left, right| left.0.cmp(right.0));
hasher.update(entries.len().to_le_bytes());
for (name, value) in entries {
Self::update_hash_str(hasher, name);
Self::update_hash_str(hasher, value);
}
}
pub(crate) fn file_info_quorum_hash(meta: &FileInfo) -> [u8; 32] {
let mut hasher = Sha256::new();
Self::update_file_info_quorum_hash(&mut hasher, meta);
let digest = hasher.finalize();
let mut key = [0u8; 32];
key.copy_from_slice(digest.as_slice());
key
}
/// Hash the per-target delete-marker versions through their normalized form
/// so the dual internal prefixes carrying the same mapping share one
/// identity, while a genuine disagreement between disks still changes the
/// hash and surfaces as a quorum difference.
fn update_hash_target_delete_marker_versions(hasher: &mut Sha256, metadata: &HashMap<String, String>) {
let (versions, corrupt) = http::target_delete_marker_versions(metadata);
hasher.update([u8::from(corrupt)]);
let mut versions = versions.iter().collect::<Vec<_>>();
versions.sort_by(|left, right| left.0.cmp(right.0));
hasher.update(versions.len().to_le_bytes());
for (arn, version_id) in versions {
Self::update_hash_str(hasher, arn);
Self::update_hash_str(hasher, version_id);
}
}
fn update_file_info_quorum_hash(hasher: &mut Sha256, meta: &FileInfo) {
hasher.update(meta.size.to_le_bytes());
hasher.update([u8::from(meta.deleted), u8::from(meta.mark_deleted)]);
hasher.update([u8::from(meta.expire_restored)]);
hasher.update([
u8::from(meta.is_remote()),
u8::from(Self::file_info_has_encryption_metadata(meta)),
u8::from(meta.is_compressed()),
]);
Self::update_hash_optional_time(hasher, meta.mod_time);
Self::update_hash_str(hasher, &meta.transition_status);
Self::update_hash_str(hasher, &meta.transition_tier);
Self::update_hash_str(hasher, &meta.transitioned_objname);
Self::update_hash_optional_uuid(hasher, meta.transition_version_id);
Self::update_hash_optional_str(hasher, meta.transition_version.as_deref());
hasher.update([match meta.transition_version_state {
rustfs_filemeta::TransitionVersionState::Unknown => 0,
rustfs_filemeta::TransitionVersionState::KnownDisabled => 1,
rustfs_filemeta::TransitionVersionState::SuspendedNull => 2,
rustfs_filemeta::TransitionVersionState::Exact => 3,
}]);
Self::update_hash_optional_u32(hasher, meta.mode);
Self::update_hash_optional_u64(hasher, meta.written_by_version);
Self::update_hash_optional_uuid(hasher, meta.version_id);
Self::update_hash_optional_uuid(hasher, meta.data_dir);
Self::update_hash_optional_bytes(hasher, meta.checksum.as_ref());
Self::update_hash_quorum_metadata_map(hasher, &meta.metadata);
Self::update_hash_target_delete_marker_versions(hasher, &meta.metadata);
hasher.update(meta.parts.len().to_le_bytes());
for part in meta.parts.iter() {
hasher.update(part.number.to_le_bytes());
hasher.update(part.size.to_le_bytes());
hasher.update(part.actual_size.to_le_bytes());
Self::update_hash_str(hasher, &part.etag);
Self::update_hash_optional_time(hasher, part.mod_time);
Self::update_hash_optional_bytes(hasher, part.index.as_ref());
Self::update_hash_optional_str(hasher, part.error.as_deref());
if let Some(checksums) = &part.checksums {
let mut checksum_entries = checksums.iter().collect::<Vec<_>>();
checksum_entries.sort_by(|left, right| left.0.cmp(right.0));
hasher.update(checksum_entries.len().to_le_bytes());
for (name, value) in checksum_entries {
Self::update_hash_str(hasher, name);
Self::update_hash_str(hasher, value);
}
} else {
hasher.update(0usize.to_le_bytes());
}
}
if !meta.is_canonical_delete_marker() && meta.size != 0 {
hasher.update(meta.erasure.data_blocks.to_le_bytes());
hasher.update(meta.erasure.parity_blocks.to_le_bytes());
hasher.update(meta.erasure.distribution.len().to_le_bytes());
for disk_index in meta.erasure.distribution.iter() {
hasher.update(disk_index.to_le_bytes());
}
}
}
fn latest_fileinfo_identity_groups(parts_metadata: &[FileInfo], errs: &[Option<DiskError>]) -> Vec<FileInfoIdentityGroup> {
let mut groups: Vec<FileInfoIdentityGroup> = Vec::with_capacity(parts_metadata.len());
for (meta, err) in parts_metadata.iter().zip(errs.iter()) {
if err.is_some() || !file_info_is_valid_for_metadata(meta) {
continue;
}
let hash = Self::file_info_quorum_hash(meta);
if let Some(group) = groups.iter_mut().find(|group| group.hash == hash) {
group.count += 1;
continue;
}
groups.push(FileInfoIdentityGroup {
hash,
count: 1,
mod_time: meta.mod_time,
});
}
groups
}
fn pick_fileinfo_identity(
disks: &[Option<DiskStore>],
parts_metadata: &[FileInfo],
errs: &[Option<DiskError>],
hash: [u8; 32],
quorum: usize,
) -> disk::error::Result<(Vec<Option<DiskStore>>, FileInfo)> {
let mut online_disks = vec![None; disks.len()];
let mut selected = None;
let mut count = 0;
for (i, ((meta, err), disk)) in parts_metadata.iter().zip(errs.iter()).zip(disks.iter()).enumerate() {
if err.is_some() || !file_info_is_valid_for_metadata(meta) || Self::file_info_quorum_hash(meta) != hash {
continue;
}
count += 1;
online_disks[i].clone_from(disk);
if selected.is_none() {
selected = Some(meta.clone());
}
}
if count < quorum {
return Err(DiskError::ErasureReadQuorum);
}
selected
.map(|mut fi| {
fi.is_latest = fi.successor_mod_time.is_none();
(online_disks, fi)
})
.ok_or(DiskError::ErasureReadQuorum)
}
fn pick_degraded_latest_fileinfo(
disks: &[Option<DiskStore>],
parts_metadata: &[FileInfo],
errs: &[Option<DiskError>],
read_quorum: usize,
write_quorum: usize,
) -> disk::error::Result<(Vec<Option<DiskStore>>, FileInfo)> {
let mut groups = Self::latest_fileinfo_identity_groups(parts_metadata, errs);
if groups.is_empty() {
return Err(DiskError::ErasureReadQuorum);
}
groups.sort_by(|left, right| right.mod_time.cmp(&left.mod_time).then_with(|| right.count.cmp(&left.count)));
let latest_mod_time = groups[0].mod_time;
let mut older_start = 0;
while older_start < groups.len() && groups[older_start].mod_time == latest_mod_time {
if groups[older_start].count >= write_quorum {
return Self::pick_fileinfo_identity(disks, parts_metadata, errs, groups[older_start].hash, write_quorum);
}
older_start += 1;
}
if older_start > 1 {
return Err(DiskError::ErasureReadQuorum);
}
for group in groups.iter().skip(older_start) {
if group.count >= read_quorum {
return Self::pick_fileinfo_identity(disks, parts_metadata, errs, group.hash, read_quorum);
}
}
Err(DiskError::ErasureReadQuorum)
}
pub(super) fn find_file_info_in_quorum(
metas: &[FileInfo],
mod_time: &Option<OffsetDateTime>,
etag: &Option<String>,
quorum: usize,
) -> disk::error::Result<FileInfo> {
if quorum < 1 {
warn!("find_file_info_in_quorum: quorum < 1");
return Err(DiskError::ErasureReadQuorum);
}
let mut meta_hashes = vec![None; metas.len()];
for (i, meta) in metas.iter().enumerate() {
if !file_info_is_valid_for_metadata(meta) {
debug!(
index = i,
valid = false,
version_id = ?meta.version_id,
mod_time = ?meta.mod_time,
"find_file_info_in_quorum: skipping invalid meta"
);
continue;
}
debug!(
index = i,
valid = true,
version_id = ?meta.version_id,
mod_time = ?meta.mod_time,
deleted = meta.deleted,
size = meta.size,
"find_file_info_in_quorum: inspecting meta"
);
let etag_only = mod_time.is_none()
&& etag.is_some()
&& meta
.get_etag()
.is_some_and(|v| &v == etag.as_ref().expect("operation should succeed"));
let mod_valid = mod_time == &meta.mod_time;
if etag_only || mod_valid {
meta_hashes[i] = Some(Self::file_info_quorum_hash(meta));
} else {
debug!(
index = i,
etag_only_match = etag_only,
mod_valid_match = mod_valid,
"find_file_info_in_quorum: meta does not match common etag or mod_time, skipping hash calculation"
);
}
}
let mut count_map = HashMap::new();
for hash in meta_hashes.iter().flatten().copied() {
*count_map.entry(hash).or_insert(0) += 1;
}
let mut max_val = None;
let mut max_count = 0;
for (&val, &count) in &count_map {
if count > max_count {
max_val = Some(val);
max_count = count;
}
}
if max_count < quorum {
warn!(
quorum,
max_count,
max_val = ?max_val,
count_map = ?count_map,
"find_file_info_in_quorum: fileinfo content identity did not reach quorum"
);
return Err(DiskError::ErasureReadQuorum);
}
let mut found_fi = None;
let mut found = false;
let mut valid_obj_map = HashMap::new();
for (i, op_hash) in meta_hashes.iter().enumerate() {
if let Some(hash) = op_hash
&& let Some(max_hash) = max_val
&& *hash == max_hash
&& file_info_is_valid_for_metadata(&metas[i])
{
if !found {
found_fi = Some(metas[i].clone());
found = true;
}
let props = ObjProps {
successor_mod_time: metas[i].successor_mod_time,
num_versions: metas[i].num_versions,
};
*valid_obj_map.entry(props).or_insert(0) += 1;
}
}
if found {
let mut fi = found_fi.expect("operation should succeed");
for (val, &count) in &valid_obj_map {
if count >= quorum {
fi.successor_mod_time = val.successor_mod_time;
fi.num_versions = val.num_versions;
fi.is_latest = val.successor_mod_time.is_none();
break;
}
}
return Ok(fi);
}
warn!("find_file_info_in_quorum: fileinfo not found");
Err(DiskError::ErasureReadQuorum)
}
/// Ownership-taking variant of `shuffle_disks_and_parts_metadata_by_index`
/// (backlog#873): callers that already own the vectors avoid one deep
/// `FileInfo` clone per disk by moving entries into their shuffled slots.
///
/// Semantics match the borrowing variant, including the fallback to the
/// mod-time based placement when `parity_blocks` or more sources are
/// inconsistent; the consistency check runs as a read-only first pass so
/// the fallback still sees the untouched inputs.
pub(super) fn shuffle_disks_and_parts_metadata_by_index_owned(
mut disks: Vec<Option<DiskStore>>,
mut parts_metadata: Vec<FileInfo>,
fi: &FileInfo,
) -> (Vec<Option<DiskStore>>, Vec<FileInfo>) {
let distribution = &fi.erasure.distribution;
let mut inconsistent = 0;
for (k, v) in parts_metadata.iter().enumerate() {
if disks[k].is_none() || !v.has_valid_erasure_geometry() || distribution[k] != v.erasure.index {
inconsistent += 1;
}
}
let use_by_index = inconsistent < fi.erasure.parity_blocks;
let init = fi.mod_time.is_none();
let mut shuffled_disks = vec![None; disks.len()];
let mut shuffled_parts_metadata = vec![FileInfo::default(); parts_metadata.len()];
for k in 0..parts_metadata.len() {
if disks[k].is_none() {
continue;
}
let eligible = if use_by_index {
parts_metadata[k].has_valid_erasure_geometry() && distribution[k] == parts_metadata[k].erasure.index
} else {
init || parts_metadata[k].has_valid_erasure_geometry()
};
if !eligible {
continue;
}
// Defensive: a corrupt/adversarial `distribution` value of `0` would
// underflow `block_idx - 1`, and a value `> N` would index out of
// bounds. Skip such entries instead of panicking (backlog#949).
let Some(slot) = distribution[k]
.checked_sub(1)
.filter(|slot| *slot < shuffled_parts_metadata.len() && *slot < shuffled_disks.len())
else {
continue;
};
shuffled_parts_metadata[slot] = std::mem::take(&mut parts_metadata[k]);
shuffled_disks[slot] = disks[k].take();
}
(shuffled_disks, shuffled_parts_metadata)
}
pub(super) fn shuffle_disks_and_parts_metadata_by_index(
disks: &[Option<DiskStore>],
parts_metadata: &[FileInfo],
fi: &FileInfo,
) -> (Vec<Option<DiskStore>>, Vec<FileInfo>) {
let mut shuffled_disks = vec![None; disks.len()];
let mut shuffled_parts_metadata = vec![FileInfo::default(); parts_metadata.len()];
let distribution = &fi.erasure.distribution;
let mut inconsistent = 0;
for (k, v) in parts_metadata.iter().enumerate() {
if disks[k].is_none() {
inconsistent += 1;
continue;
}
if !v.has_valid_erasure_geometry() {
inconsistent += 1;
continue;
}
if distribution[k] != v.erasure.index {
inconsistent += 1;
continue;
}
// Defensive: reject out-of-range distribution values instead of
// underflowing/indexing out of bounds (backlog#949).
let Some(slot) = distribution[k]
.checked_sub(1)
.filter(|slot| *slot < shuffled_parts_metadata.len() && *slot < shuffled_disks.len())
else {
inconsistent += 1;
continue;
};
shuffled_parts_metadata[slot] = parts_metadata[k].clone();
shuffled_disks[slot].clone_from(&disks[k]);
}
if inconsistent < fi.erasure.parity_blocks {
return (shuffled_disks, shuffled_parts_metadata);
}
Self::shuffle_disks_and_parts_metadata(disks, parts_metadata, fi)
}
pub(super) fn shuffle_disks_and_parts_metadata(
disks: &[Option<DiskStore>],
parts_metadata: &[FileInfo],
fi: &FileInfo,
) -> (Vec<Option<DiskStore>>, Vec<FileInfo>) {
let init = fi.mod_time.is_none();
let mut shuffled_disks = vec![None; disks.len()];
let mut shuffled_parts_metadata = vec![FileInfo::default(); parts_metadata.len()];
let distribution = &fi.erasure.distribution;
for (k, v) in disks.iter().enumerate() {
if v.is_none() {
continue;
}
if !init && !parts_metadata[k].has_valid_erasure_geometry() {
continue;
}
// if !init && fi.xlv1 != parts_metadata[k].xlv1 {
// continue;
// }
// Defensive: reject out-of-range distribution values instead of
// underflowing/indexing out of bounds (backlog#949).
let Some(slot) = distribution[k]
.checked_sub(1)
.filter(|slot| *slot < shuffled_parts_metadata.len() && *slot < shuffled_disks.len())
else {
continue;
};
shuffled_parts_metadata[slot] = parts_metadata[k].clone();
shuffled_disks[slot].clone_from(&disks[k]);
}
(shuffled_disks, shuffled_parts_metadata)
}
pub(super) fn shuffle_parts_metadata(parts_metadata: &[FileInfo], distribution: &[usize]) -> Vec<FileInfo> {
if distribution.is_empty() {
return parts_metadata.to_vec();
}
let mut shuffled_parts_metadata = vec![FileInfo::default(); parts_metadata.len()];
// Shuffle slice xl metadata for expected distribution.
for (index, part) in parts_metadata.iter().enumerate() {
// Defensive: skip missing or out-of-range distribution values
// instead of underflowing/indexing out of bounds (backlog#949).
let Some(slot) = distribution
.get(index)
.and_then(|block_index| block_index.checked_sub(1))
.filter(|slot| *slot < shuffled_parts_metadata.len())
else {
continue;
};
shuffled_parts_metadata[slot] = part.clone();
}
shuffled_parts_metadata
}
pub(super) fn shuffle_disks(disks: &[Option<DiskStore>], distribution: &[usize]) -> Vec<Option<DiskStore>> {
if distribution.is_empty() {
return disks.to_vec();
}
let mut shuffled_disks = vec![None; disks.len()];
for (i, v) in disks.iter().enumerate() {
// Defensive: skip missing or out-of-range distribution values
// instead of underflowing/indexing out of bounds (backlog#949).
let Some(slot) = distribution
.get(i)
.and_then(|idx| idx.checked_sub(1))
.filter(|slot| *slot < shuffled_disks.len())
else {
continue;
};
shuffled_disks[slot].clone_from(v);
}
shuffled_disks
}
pub(super) fn shuffle_disks_owned(mut disks: Vec<Option<DiskStore>>, distribution: &[usize]) -> Vec<Option<DiskStore>> {
if distribution.is_empty() {
return disks;
}
let mut shuffled_disks = vec![None; disks.len()];
for (index, disk) in disks.iter_mut().enumerate() {
let Some(slot) = distribution
.get(index)
.and_then(|block_index| block_index.checked_sub(1))
.filter(|slot| *slot < shuffled_disks.len())
else {
continue;
};
shuffled_disks[slot] = disk.take();
}
shuffled_disks
}
pub(super) fn shuffle_check_parts(parts_errs: &[usize], distribution: &[usize]) -> Vec<usize> {
if distribution.is_empty() {
return parts_errs.to_vec();
}
let mut shuffled_parts_errs = vec![0; parts_errs.len()];
for (i, v) in parts_errs.iter().enumerate() {
// Defensive: skip missing or out-of-range distribution values
// instead of underflowing/indexing out of bounds (backlog#949).
let Some(slot) = distribution
.get(i)
.and_then(|idx| idx.checked_sub(1))
.filter(|slot| *slot < shuffled_parts_errs.len())
else {
continue;
};
shuffled_parts_errs[slot] = *v;
}
shuffled_parts_errs
}
}
#[cfg(test)]
mod tests {
use super::*;
fn metadata_quorum_test_fileinfo(mod_time: OffsetDateTime, erasure_index: usize) -> FileInfo {
let mut fi = FileInfo::new("bucket/object", 2, 2);
fi.name = "bucket/object".to_string();
fi.size = 8 * 1024 * 1024;
fi.mod_time = Some(mod_time);
fi.data_dir = Some(Uuid::new_v4());
fi.metadata.insert("etag".to_string(), "object-etag".to_string());
fi.add_object_part(1, "part-etag".to_string(), 8 * 1024 * 1024, Some(mod_time), 8 * 1024 * 1024, None, None);
fi.erasure.index = erasure_index;
fi
}
fn transition_metadata_quorum_fileinfo(erasure_index: usize) -> FileInfo {
let mut fi = FileInfo::new("bucket/object", 5, 1);
fi.name = "bucket/object".to_string();
fi.size = 8 * 1024 * 1024;
fi.mod_time = Some(OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"));
fi.metadata.insert("etag".to_string(), "object-etag".to_string());
fi.transition_status = TRANSITION_COMPLETE.to_string();
fi.transition_tier = "WARM".to_string();
fi.transitioned_objname = "remote/object".to_string();
fi.transition_version_id = Some(Uuid::new_v4());
fi.erasure.index = erasure_index;
fi
}
fn expect_metadata_quorum_error(metas: Vec<FileInfo>, mod_time: OffsetDateTime, message: &str) {
let err = SetDisks::find_file_info_in_quorum(&metas, &Some(mod_time), &None, 3).expect_err(message);
assert_eq!(err, DiskError::ErasureReadQuorum);
}
#[test]
fn metadata_quorum_covers_etag_fallback_and_object_quorum_failures() {
let mut parts_metadata = (1..=3)
.map(|index| {
let mut fi = metadata_quorum_test_fileinfo(OffsetDateTime::now_utc(), index);
fi.mod_time = None;
fi
})
.collect::<Vec<_>>();
parts_metadata[2]
.metadata
.insert("etag".to_string(), "minority-etag".to_string());
let errs = vec![None; parts_metadata.len()];
let disks = vec![None; parts_metadata.len()];
let (_online, mod_time, etag) = SetDisks::list_online_disks(&disks, &parts_metadata, &errs, 2);
assert!(mod_time.is_none());
assert_eq!(etag.as_deref(), Some("object-etag"));
let zero_parity = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 0)
.expect("zero default parity should require all metadata shards");
assert_eq!(zero_parity, (3, 3));
let invalid = vec![FileInfo::default(); 4];
let err = SetDisks::object_quorum_from_meta(&invalid, &vec![None; 4], 2)
.expect_err("invalid metadata without a common parity must fail closed");
// A full set of healthy replies whose metadata cannot be reconciled is
// corrupt (heal-actionable), not a retryable quorum outage (rustfs#5801).
assert_eq!(err, DiskError::FileCorrupt);
}
#[test]
fn fileinfo_quorum_hash_includes_optional_checksums_and_ignores_replication_noise() {
let mod_time = OffsetDateTime::now_utc();
let mut left = metadata_quorum_test_fileinfo(mod_time, 1);
left.mode = Some(0o640);
left.written_by_version = Some(42);
left.checksum = Some(Bytes::from_static(b"object-checksum"));
left.parts[0].index = Some(Bytes::from_static(b"part-index"));
left.parts[0].error = Some("repair-pending".to_string());
left.parts[0].checksums = Some(HashMap::from([
("sha256".to_string(), "left".to_string()),
("crc32".to_string(), "right".to_string()),
]));
left.metadata.insert(
format!("{}{}", http::RUSTFS_INTERNAL_PREFIX, http::SUFFIX_REPLICATION_STATUS),
"replica-a".to_string(),
);
let mut right = left.clone();
right.parts[0].checksums = Some(HashMap::from([
("crc32".to_string(), "right".to_string()),
("sha256".to_string(), "left".to_string()),
]));
right.metadata.insert(
format!("{}{}", http::MINIO_INTERNAL_PREFIX, http::SUFFIX_REPLICATION_STATUS),
"replica-b".to_string(),
);
assert_eq!(
SetDisks::file_info_quorum_hash(&left),
SetDisks::file_info_quorum_hash(&right),
"checksum map ordering and replication metadata must not split quorum identity"
);
right.parts[0]
.checksums
.as_mut()
.expect("checksums should exist")
.insert("sha256".to_string(), "changed".to_string());
assert_ne!(SetDisks::file_info_quorum_hash(&left), SetDisks::file_info_quorum_hash(&right));
}
#[test]
fn purge_pending_quorum_hash_keeps_erasure_layouts_separate() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let version_id = Uuid::new_v4();
let data_dir = Uuid::new_v4();
let mut honest = FileInfo::new("bucket/object", 5, 1);
honest.name = "bucket/object".to_string();
honest.version_id = Some(version_id);
honest.data_dir = Some(data_dir);
honest.mod_time = Some(mod_time);
honest.size = 1;
honest.deleted = true;
honest.add_object_part(1, "part-etag".to_string(), 1, Some(mod_time), 1, None, None);
let mut parts_metadata = (1..=6)
.map(|index| {
let mut metadata = honest.clone();
metadata.erasure.index = index;
metadata
})
.collect::<Vec<_>>();
let mut tampered_layout = FileInfo::new("bucket/object", 3, 3).erasure;
tampered_layout.index = 1;
parts_metadata[0].erasure = tampered_layout;
let errs = vec![None; 6];
assert_eq!(
SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 3)
.expect("five honest EC:1 payload copies should determine object quorum"),
(5, 5)
);
let selected = SetDisks::find_file_info_in_quorum(&parts_metadata, &Some(mod_time), &None, 5)
.expect("the five matching EC:1 payload copies should determine metadata identity");
assert_eq!(selected.erasure.data_blocks, 5);
assert_eq!(selected.erasure.parity_blocks, 1);
}
#[test]
fn quorum_helpers_reject_zero_quorum_and_shuffle_check_parts_by_distribution() {
let err = SetDisks::find_file_info_in_quorum(&[], &None, &None, 0).expect_err("zero quorum cannot select metadata");
assert_eq!(err, DiskError::ErasureReadQuorum);
assert_eq!(SetDisks::shuffle_check_parts(&[2, 1, 0], &[]), vec![2, 1, 0]);
assert_eq!(SetDisks::shuffle_check_parts(&[2, 1, 0], &[3, 1, 2]), vec![1, 0, 2]);
}
#[test]
fn metadata_quorum_uses_simple_majority_for_transitioned_objects() {
let parts_metadata = (1..=6).map(transition_metadata_quorum_fileinfo).collect::<Vec<_>>();
let errs = vec![None; parts_metadata.len()];
let parities = SetDisks::list_object_parities(&parts_metadata, &errs);
let quorum = SetDisks::object_quorum_from_meta(&parts_metadata, &errs, 1)
.expect("transitioned metadata should resolve object quorum");
assert_eq!(parities, vec![2; 6]);
assert_eq!(quorum, (4, 4));
}
#[test]
fn find_file_info_in_quorum_rejects_encrypted_plain_metadata_split() {
let mod_time = OffsetDateTime::now_utc();
let mut encrypted_a = metadata_quorum_test_fileinfo(mod_time, 1);
let mut encrypted_b = metadata_quorum_test_fileinfo(mod_time, 2);
let plain = metadata_quorum_test_fileinfo(mod_time, 3);
encrypted_a
.metadata
.insert("x-rustfs-encryption-key".to_string(), "encrypted-key".to_string());
encrypted_b
.metadata
.insert("x-rustfs-encryption-key".to_string(), "encrypted-key".to_string());
expect_metadata_quorum_error(
vec![encrypted_a, encrypted_b, plain],
mod_time,
"mixed encrypted and plain metadata must not reach quorum",
);
}
#[test]
fn find_file_info_in_quorum_rejects_compressed_plain_metadata_split() {
let mod_time = OffsetDateTime::now_utc();
let mut compressed_a = metadata_quorum_test_fileinfo(mod_time, 1);
let mut compressed_b = metadata_quorum_test_fileinfo(mod_time, 2);
let plain = metadata_quorum_test_fileinfo(mod_time, 3);
http::insert_str(&mut compressed_a.metadata, http::SUFFIX_COMPRESSION, "lz4".to_string());
http::insert_str(&mut compressed_b.metadata, http::SUFFIX_COMPRESSION, "lz4".to_string());
expect_metadata_quorum_error(
vec![compressed_a, compressed_b, plain],
mod_time,
"mixed compressed and plain metadata must not reach quorum",
);
}
#[test]
fn find_file_info_in_quorum_rejects_remote_local_metadata_split() {
let mod_time = OffsetDateTime::now_utc();
let mut remote = metadata_quorum_test_fileinfo(mod_time, 1);
let local_a = metadata_quorum_test_fileinfo(mod_time, 2);
let local_b = metadata_quorum_test_fileinfo(mod_time, 3);
remote.transition_status = TRANSITION_COMPLETE.to_string();
remote.transition_tier = "WARM".to_string();
remote.transitioned_objname = "remote/object".to_string();
remote.transition_version_id = Some(Uuid::new_v4());
expect_metadata_quorum_error(
vec![remote, local_a, local_b],
mod_time,
"mixed remote and local metadata must not reach quorum",
);
}
async fn shuffle_test_disks(tempdir: &tempfile::TempDir, count: usize) -> Vec<Option<DiskStore>> {
let endpoint =
Endpoint::try_from(tempdir.path().to_str().expect("tempdir path should be utf8")).expect("endpoint should parse");
let disk = new_disk(
&endpoint,
&DiskOption {
cleanup: false,
health_check: false,
},
)
.await
.expect("disk should be created");
// The shuffle only inspects Some/None and clones the Arc handle, so
// one shared disk handle per slot is sufficient.
(0..count).map(|_| Some(disk.clone())).collect()
}
fn shuffle_fixture(consistent: bool) -> (FileInfo, Vec<FileInfo>) {
let mut fi = FileInfo::new("bucket/object", 2, 1);
fi.mod_time = Some(time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp"));
fi.size = 1;
fi.add_object_part(1, String::new(), 1, None, 1, None, None);
let slots = fi.erasure.distribution.len();
let parts = (0..slots)
.map(|k| {
let mut part_fi = fi.clone();
part_fi.erasure.index = if consistent {
fi.erasure.distribution[k]
} else {
// Misplace every source so the by-index pass is rejected
// and the mod-time fallback placement runs instead.
fi.erasure.distribution[(k + 1) % slots]
};
part_fi
})
.collect();
(fi, parts)
}
#[tokio::test]
async fn owned_shuffle_matches_borrowing_variant_when_consistent() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let (fi, parts) = shuffle_fixture(true);
let disks = shuffle_test_disks(&tempdir, parts.len()).await;
let (expected_disks, expected_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index(&disks, &parts, &fi);
let (owned_disks, owned_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index_owned(disks, parts, &fi);
assert_eq!(owned_parts, expected_parts, "owned shuffle must place identical metadata");
let expected_slots: Vec<bool> = expected_disks.iter().map(Option::is_some).collect();
let owned_slots: Vec<bool> = owned_disks.iter().map(Option::is_some).collect();
assert_eq!(owned_slots, expected_slots, "owned shuffle must fill identical disk slots");
}
#[tokio::test]
async fn owned_shuffle_matches_borrowing_variant_on_fallback() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let (fi, parts) = shuffle_fixture(false);
let disks = shuffle_test_disks(&tempdir, parts.len()).await;
let (expected_disks, expected_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index(&disks, &parts, &fi);
let (owned_disks, owned_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index_owned(disks, parts, &fi);
assert_eq!(owned_parts, expected_parts, "fallback placement must match the borrowing variant");
let expected_slots: Vec<bool> = expected_disks.iter().map(Option::is_some).collect();
let owned_slots: Vec<bool> = owned_disks.iter().map(Option::is_some).collect();
assert_eq!(owned_slots, expected_slots, "fallback disk slots must match the borrowing variant");
}
#[tokio::test]
async fn owned_shuffle_preserves_fresh_put_metadata() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let fi = FileInfo::new("bucket/object", 2, 1);
let parts = vec![fi.clone(); fi.erasure.distribution.len()];
let disks = shuffle_test_disks(&tempdir, parts.len()).await;
let (owned_disks, owned_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index_owned(disks, parts, &fi);
assert!(owned_disks.iter().all(Option::is_some), "fresh PUT must retain every online disk");
assert_eq!(
owned_parts,
vec![fi; owned_disks.len()],
"fresh PUT metadata with pending shard indexes must survive init fallback"
);
}
// backlog#949: corrupt/adversarial distribution values (0 or > N) must not
// trigger a `usize` underflow / out-of-bounds panic in the shuffle helpers.
#[test]
fn shuffle_parts_metadata_survives_corrupt_distribution() {
let parts = vec![FileInfo::default(); 4];
// distribution[0] = 0 underflows; distribution[1] = 9 is out of range.
let result = SetDisks::shuffle_parts_metadata(&parts, &[0, 9, 3, 4]);
assert_eq!(result.len(), parts.len(), "output length must be preserved");
}
#[test]
fn shuffle_check_parts_survives_corrupt_distribution() {
let errs = vec![10usize, 20, 30, 40];
let result = SetDisks::shuffle_check_parts(&errs, &[0, 9, 3, 4]);
assert_eq!(result.len(), errs.len(), "output length must be preserved");
// In-range entries are still placed; corrupt ones are safely skipped.
assert_eq!(result[2], 30, "distribution[2]=3 places errs[2] into slot 2");
assert_eq!(result[3], 40, "distribution[3]=4 places errs[3] into slot 3");
}
#[tokio::test]
async fn shuffle_disks_survives_corrupt_distribution() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let disks = shuffle_test_disks(&tempdir, 4).await;
// distribution[0] = 0 underflows; distribution[1] = 9 is out of range.
let result = SetDisks::shuffle_disks(&disks, &[0, 9, 3, 4]);
assert_eq!(result.len(), disks.len(), "output length must be preserved");
}
#[tokio::test]
async fn owned_disk_shuffle_matches_borrowing_variant() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let mut disks = shuffle_test_disks(&tempdir, 4).await;
disks[1] = None;
disks[3] = None;
let distribution = [3, 1, 4, 2];
let expected = SetDisks::shuffle_disks(&disks, &distribution);
let actual = SetDisks::shuffle_disks_owned(disks, &distribution);
let expected_slots = expected.iter().map(Option::is_some).collect::<Vec<_>>();
let actual_slots = actual.iter().map(Option::is_some).collect::<Vec<_>>();
assert_eq!(actual_slots, expected_slots, "owned shuffle must preserve disk placement");
}
#[tokio::test]
async fn shuffle_disks_and_parts_metadata_survives_corrupt_distribution() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let (mut fi, parts) = shuffle_fixture(true);
let slots = parts.len();
// Corrupt the selected FileInfo's distribution: a `0` (underflow) and an
// out-of-range value. Length and `erasure.index` stay well-formed so the
// corruption is only in the distribution values.
fi.erasure.distribution = vec![0; slots];
fi.erasure.distribution[0] = slots + 5;
let disks = shuffle_test_disks(&tempdir, slots).await;
// None of these must panic on the corrupt distribution.
let (d1, _) = SetDisks::shuffle_disks_and_parts_metadata(&disks, &parts, &fi);
assert_eq!(d1.len(), disks.len());
let (d2, _) = SetDisks::shuffle_disks_and_parts_metadata_by_index(&disks, &parts, &fi);
assert_eq!(d2.len(), disks.len());
let (d3, _) = SetDisks::shuffle_disks_and_parts_metadata_by_index_owned(disks, parts, &fi);
assert_eq!(d3.len(), slots);
}
#[tokio::test]
async fn shuffle_variants_skip_missing_disks_and_invalid_metadata() {
let tempdir = tempfile::tempdir().expect("tempdir should be created");
let (fi, mut parts) = shuffle_fixture(true);
let mut disks = shuffle_test_disks(&tempdir, parts.len()).await;
disks[0] = None;
parts[1] = FileInfo::default();
let (by_index_disks, by_index_parts) = SetDisks::shuffle_disks_and_parts_metadata_by_index(&disks, &parts, &fi);
assert!(
by_index_disks.iter().filter(|disk| disk.is_some()).count() <= disks.iter().filter(|disk| disk.is_some()).count()
);
assert!(by_index_parts.iter().any(|part| !part.is_valid()));
let (fallback_disks, fallback_parts) = SetDisks::shuffle_disks_and_parts_metadata(&disks, &parts, &fi);
assert!(fallback_disks.iter().any(Option::is_none));
assert!(fallback_parts.iter().any(|part| !part.is_valid()));
}
#[test]
fn target_delete_marker_version_metadata_is_included_in_quorum_hash() {
let suffix = "replication-delete-marker-version-arn:rustfs:replication::target:bucket";
assert!(SetDisks::is_replication_quorum_metadata_key(&format!(
"{}{}",
http::RUSTFS_INTERNAL_PREFIX,
suffix
)));
assert!(SetDisks::is_replication_quorum_metadata_key(&format!(
"{}{}",
http::MINIO_INTERNAL_PREFIX,
suffix
)));
assert!(!SetDisks::is_replication_quorum_metadata_key("x-rustfs-internal-unrelated"));
let mut left = metadata_quorum_test_fileinfo(OffsetDateTime::now_utc(), 1);
let mut right = left.clone();
left.metadata
.insert(format!("{}{}", http::RUSTFS_INTERNAL_PREFIX, suffix), "target-version-a".to_string());
right
.metadata
.insert(format!("{}{}", http::MINIO_INTERNAL_PREFIX, suffix), "target-version-b".to_string());
assert_ne!(SetDisks::file_info_quorum_hash(&left), SetDisks::file_info_quorum_hash(&right));
let mut dual_prefixed = left.clone();
dual_prefixed
.metadata
.insert(format!("{}{}", http::MINIO_INTERNAL_PREFIX, suffix), "target-version-a".to_string());
assert_eq!(
SetDisks::file_info_quorum_hash(&left),
SetDisks::file_info_quorum_hash(&dual_prefixed),
"compatible prefixes carrying the same mapping must share one identity"
);
}
/// Guards the switch to `rustfs_utils::http::starts_with_ignore_ascii_case`:
/// internal prefixes must keep matching case-insensitively, and keys shorter
/// than the prefix must keep being rejected. Misclassifying either way leaks
/// internal metadata into the quorum hash (or drops it out of it).
#[test]
fn internal_metadata_suffix_is_prefix_case_insensitive_and_rejects_short_keys() {
assert_eq!(
SetDisks::internal_metadata_suffix("X-RustFS-Internal-Replica-Status"),
Some("Replica-Status"),
"mixed-case RustFS prefix must match and preserve the suffix casing"
);
assert_eq!(
SetDisks::internal_metadata_suffix("X-MINIO-INTERNAL-replica-status"),
Some("replica-status"),
"mixed-case MinIO prefix must match"
);
assert_eq!(SetDisks::internal_metadata_suffix(http::RUSTFS_INTERNAL_PREFIX), Some(""));
// Keys shorter than either prefix, and non-internal keys, stay unmatched.
assert_eq!(SetDisks::internal_metadata_suffix(""), None);
assert_eq!(SetDisks::internal_metadata_suffix("x-rustfs-interna"), None);
assert_eq!(SetDisks::internal_metadata_suffix("x-minio-interna"), None);
assert_eq!(SetDisks::internal_metadata_suffix("x-amz-meta-custom"), None);
// The suffix-prefix comparisons behind the classifier follow the same rules.
assert!(SetDisks::is_replication_quorum_metadata_key(
"X-RustFS-Internal-Replication-Reset-arn:rustfs:replication::target:bucket"
));
assert!(SetDisks::is_replication_quorum_metadata_key(
"X-Minio-Internal-Replication-Delete-Marker-Version-arn:rustfs:replication::target:bucket"
));
assert!(!SetDisks::is_replication_quorum_metadata_key("x-rustfs-interna"));
assert!(!SetDisks::is_replication_quorum_metadata_key("x-rustfs-internal-replication-res"));
}
/// rustfs#5801: parity counts outside [0, total_shards] come from corrupt
/// or foreign metadata and must be treated as invalid entries instead of
/// clamped values that poison `common_parity`'s occurrence counting.
#[test]
fn out_of_range_parity_is_treated_as_invalid_entry() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut metas: Vec<FileInfo> = (0..4).map(|i| metadata_quorum_test_fileinfo(mod_time, i)).collect();
for fi in &mut metas {
fi.erasure.parity_blocks = usize::MAX;
}
let errs: Vec<Option<DiskError>> = vec![None; 4];
let parities = SetDisks::list_object_parities(&metas, &errs);
assert_eq!(parities, vec![-1; 4], "garbage parity must not survive as a candidate");
}
/// rustfs#5801: when a read quorum of healthy disks answers but their
/// parity values are irreconcilable, the object metadata is corrupt —
/// return `FileCorrupt` (heal-actionable, non-retryable) instead of the
/// retryable-looking `ErasureReadQuorum`.
#[test]
fn irreconcilable_parity_with_healthy_quorum_is_file_corrupt() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut metas: Vec<FileInfo> = (0..4).map(|i| metadata_quorum_test_fileinfo(mod_time, i)).collect();
for fi in &mut metas {
fi.erasure.parity_blocks = usize::MAX;
}
let errs: Vec<Option<DiskError>> = vec![None; 4];
let err = SetDisks::object_quorum_from_meta(&metas, &errs, 2).expect_err("garbage parity cannot form a quorum");
assert_eq!(err, DiskError::FileCorrupt);
}
/// Too few healthy replies remains a genuine quorum condition where a
/// retry may succeed once disks recover.
#[test]
fn insufficient_healthy_replies_stays_erasure_read_quorum() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut metas: Vec<FileInfo> = (0..4).map(|i| metadata_quorum_test_fileinfo(mod_time, i)).collect();
metas[0].erasure.parity_blocks = usize::MAX;
let errs: Vec<Option<DiskError>> = vec![
None,
Some(DiskError::DiskNotFound),
Some(DiskError::DiskNotFound),
Some(DiskError::DiskNotFound),
];
let err = SetDisks::object_quorum_from_meta(&metas, &errs, 2).expect_err("one healthy reply is below quorum");
assert_eq!(err, DiskError::ErasureReadQuorum);
}
}