mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-07 13:53:12 +00:00
5237a4465d
* feat(replication): purge delete markers by the target's own version id When a delete marker is replicated, the target assigns it a version id. The purge that follows derived one from the *source* uuid instead, which is only correct when the target mirrors source version ids. A generic S3 target does not: the derived id addresses a version that does not exist there, so the purge is a no-op and the replica keeps a marker the source has already removed. Same failure class as #4401. Record the id the target reports and address it directly on purge. Data path, all of it driven by the object's internal metadata rather than the `ReplicationState` wire form, which encodes positionally and cannot carry a map: - `rustfs-utils`: the `replication-delete-marker-version-<arn>` key family, plus `strip_internal_prefix_preserving_case` — ARNs are case-sensitive and the existing `strip_internal_prefix` lowercases. - `ReplicationState` gains the map and a `..._corrupt` flag, both `#[serde(skip)]`; `ReplicatedTargetInfo` carries the per-target id. - `persist_target_delete_marker_versions` is merge-only. A delete arriving over internode RPC has an empty map, so treating it as authoritative would let a remote disk erase an id the local disk still holds. - `delete_object_version` copies the map into `fi.metadata` before dispatch, so the durable carrier crosses the wire even though the field does not. - The keys are folded into the quorum hash through their normalized form: the dual internal prefixes carrying one mapping share an identity, while a genuine disagreement between disks still shows up as a quorum difference. - `corrupt` (the prefixes disagreed) fails closed: skip the purge and warn rather than guess an id and risk destroying a live version on the target. Ported from the rc.1 branch, which cannot merge as a whole: its MRF replay rewrite collides with #5659/#5671/#5672/#5673 and regressed `MRF_PENDING_CAP`. main's MRF machinery is kept; only this capability moves across. It touches no MRF code. Two things did not survive the port, deliberately. The branch's `missing_is_complete` purge regression does not exist here — it came from its own HEAD-precheck rewrite, and main's simpler path never had it. And the branch's `MrfReplicateEntry` ordering fields are MRF-redesign scope, left behind. Verification: cargo fmt --all --check, git diff --check, cargo check --workspace --all-targets, and the suites for the four touched crates — 4070 tests, 2 pre-existing failures unrelated to this change (`system_resolver_negative_result_reaches_the_dns_allowlist`, `test_resolve_domain_preserves_system_resolver_error_provenance`; both are the sandbox DNS interception, they fail on a clean checkout too). * fix(replication): keep the layer guard happy scripts/check_architecture_migration_rules.sh matches on text, so the doc comments naming `rustfs_filemeta::` read as a cross-layer dependency even though nothing imports it. Reword them; the guard passes. * fix(replication): make the target-version cap deterministic Two defects in this PR, both found in review. The cap was applied while iterating a `HashMap`, so *which* 1000 entries survived depended on iteration order. Two disks decoding the same oversized metadata could keep different subsets, hash differently, and lose quorum — instead of both reporting the same corruption. Collect first, then truncate in `BTreeMap` order, which is total and identical everywhere. And `persist_target_delete_marker_versions` discarded the `corrupt` flag from the RPC carrier, committing a delete-marker update that looked clean while the exact remote marker identity was unknown. It now declines to merge a corrupt carrier. Because the helper only ever inserts, declining leaves the durable keys already on the object untouched, which is strictly safer than writing a mapping we cannot trust. Residual, stated rather than papered over: corruption confined to the RPC carrier is not persisted as a sentinel, so a later reader of an object that carried no durable keys still sees "legacy, no mapping" rather than "corrupt". Persisting that would need a wire-format addition; the consumer already fails closed on any corruption it can observe. New test: `target_delete_marker_versions_cap_is_deterministic_across_decodes` decodes the same 1050-entry map twice and asserts both the corrupt flag and the retained subset agree. * fix(replication): preserve multipart source mtime (#5669) * fix(kms): repair unopenable ciphertext and cover the Vault backends (#5668) * Add black-box behavior tests for KMS resilience and serialization * fix(kms): repair unopenable ciphertext across backends Black-box testing of the KMS crate surfaced several defects that make encrypted data permanently unreadable. Symmetric envelopes. The Local and Vault Transit backends returned raw cipher output from `encrypt` while `decrypt` parsed a JSON envelope, so anything sealed through the master-key path could never be opened again. Local also discarded the AES-GCM nonce. Both now emit the same envelope `decrypt` consumes, matching the Static backend. Deterministic AAD. The object layer derived AEAD additional data by serializing a `HashMap` directly. Iteration order differs per instance, so a context rebuilt from storage produced different AAD bytes than the one used to seal and the object stopped opening. Ordering by key removes that dependency, matching the Static backend's existing `context_aad`. Objects written with the default single-key context are unaffected, since a one-entry map has only one serialization. Cipher in the header projection. `metadata_to_headers` recorded the SSE mode (`AES256` / `aws:kms`), which cannot represent ChaCha20-Poly1305, so a ChaCha-sealed object came back claiming `aws:kms` and was opened with the wrong cipher. The cipher now travels in `x-rustfs-encryption-algorithm` — the header the storage layer already reads but nothing ever wrote. Objects without it fall back as before. Also: the Static backend ignored `key_spec` and always issued 256-bit data keys; Local `list_keys` hardcoded `truncated: false`, ignored `marker`, and paginated over unordered `read_dir`, so a paginating client silently saw a partial key list; and Local and Vault KV2 reported `key_id: "unknown"` from `decrypt` despite the envelope naming the master key. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * test(kms): cover both Vault backends and key rotation The behavior suite ran only against Local and Static, and its own harness documented the gap: the Vault backends had no business-capability coverage at all. Setting `RUSTFS_KMS_VAULT_TOKEN` now adds Vault KV2 and Vault Transit to every `for_each_backend` spec against a live server. That lane is what surfaced the Transit envelope defect fixed in the previous commit. `rotate` and `versioning` are advertised only by the Vault backends, so until now every capability-gated branch for them took the `UnsupportedCapability` side and the working half was never asserted — a rotation that dropped prior key versions would have gone green. The new `behavior_rotation.rs` pins that half: material sealed before a rotation still opens after it, repeated rotations accumulate versions rather than overwriting a single spare, and the history survives a restart. Two test defects fixed. `objects_round_trip_across_sizes_and_algorithms` asserted a 1-byte object differs from its own ciphertext, which collides once every 256 runs; the assertion now applies only where a collision is not realistic, and small objects stay covered by the tag check and the decrypt round-trip. `test_from_env_selects_token_file` depended on `RUSTFS_KMS_VAULT_TOKEN` being absent from the caller's environment and now clears it explicitly. The snapshots directory was also removed from `.gitignore`: insta snapshots are the assertions themselves, so leaving them untracked gives CI nothing to compare against. Only `.snap.new` scratch files are ignored now. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * test(kms): adapt behavior suite to current key APIs Rebasing onto main brought four API changes the suite predates. `DeleteKeyRequest` gained `confirm_key_id`, and immediate deletion is now gated on the server's `allow_immediate_deletion`. Scheduled deletions pass `None`; the four specs that destroy a key outright echo the key id back and opt the harness config in, which is what the gate asks of a real caller. `LocalBackupExportRequest` gained `sanitized_config`. These specs cover the key-material path, so they seal no configuration and pass `None`. `KmsCacheStats` became a named struct with real hit, miss, and eviction counters. `cache_stats_returns_an_entry_count_and_no_hit_or_miss_data` existed to pin the old placeholder behavior — that the second tuple element was always zero — which main has since fixed, so it is now `cache_stats_reports_hits_and_misses_separately` and asserts the counters actually move. Starting the service provisions the reserved probe key, so it shows up in listings and backup bundles. Exact-set assertions filter it through a new `without_probe_key` helper rather than naming it, keeping those specs about the keys they seeded. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * fix(kms): bind the AAD to the stored context bytes Review caught that canonicalizing the AAD on decrypt breaks objects sealed before canonicalization existed, and it was right. The AAD is the *serialization* of the encryption context, and `x-rustfs-encryption-context` stores that exact byte sequence: `encrypt_object` fed one `HashMap` to the AEAD and then moved the same map into the metadata the header is written from, so the stored string is byte-identical to the AAD the object was sealed under. Those objects are therefore recoverable — but only while nothing round-trips the value through a `HashMap` and re-serializes it. Recomputing sorted AAD on decrypt would have turned a readable object into a permanently unreadable one. The previous behavior was worse than the first analysis credited: it did not merely fail intermittently, it made the failure deterministic. `EncryptionMetadata` now carries `context_aad`, the bytes the object was actually sealed with. Encryption records what it fed the AEAD, the header projection stores those bytes verbatim (and preserves a legacy ordering across a re-projection rather than rewriting it into sorted form), and `headers_to_metadata` carries the stored string through untouched. Both decrypt paths, SSE-KMS and SSE-C, prefer it and fall back to canonical serialization only when no stored serialization exists. Canonicalization still applies to everything newly sealed, so the original ordering bug cannot recur. Two tests pin this: a legacy record whose sealed bytes are non-canonical must survive a full header round trip unchanged, and a context header rewritten to an equivalent-but-reordered serialization must fail authentication rather than silently re-deriving a working AAD. Both were mutation-checked against the reinstated bug on each side. Also from review: the lifecycle churn test asserted only that every request was accounted for, which holds whether the state gate exists or not, so both branches are now pinned deterministically after the churn (asserting `refused > 0` on the concurrent phase would only trade the hole for a scheduling flake). And the Local and Vault KV2 envelopes compare `encryption_context` without authenticating it — `DekCrypto` seals only the plaintext — which is now documented at both sites; closing it needs a versioned envelope, since existing ciphertext was sealed without AAD. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: ccccpj <ccccpj@outlook.com> Co-authored-by: 唐小鸭 <tangtang1251@qq.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
1472 lines
56 KiB
Rust
1472 lines
56 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::*;
|
|
use rustfs_utils::http;
|
|
|
|
#[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_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;
|
|
}
|
|
|
|
if metadata.is_canonical_delete_marker() || metadata.size == 0 {
|
|
parities[index] = half;
|
|
} else if metadata.transition_status == TRANSITION_COMPLETE {
|
|
let majority_metadata_parity = total_shards_i32 - (half + 1);
|
|
let erasure_parity = i32::try_from(metadata.erasure.parity_blocks).unwrap_or(i32::MAX);
|
|
parities[index] = majority_metadata_parity.max(erasure_parity);
|
|
} else {
|
|
parities[index] = i32::try_from(metadata.erasure.parity_blocks).unwrap_or(i32::MAX);
|
|
}
|
|
}
|
|
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 {
|
|
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)
|
|
}
|
|
|
|
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 starts_with_ignore_ascii_case(value: &str, prefix: &str) -> bool {
|
|
value
|
|
.get(..prefix.len())
|
|
.is_some_and(|value_prefix| value_prefix.eq_ignore_ascii_case(prefix))
|
|
}
|
|
|
|
fn internal_metadata_suffix(name: &str) -> Option<&str> {
|
|
name.get(http::RUSTFS_INTERNAL_PREFIX.len()..)
|
|
.filter(|_| Self::starts_with_ignore_ascii_case(name, http::RUSTFS_INTERNAL_PREFIX))
|
|
.or_else(|| {
|
|
name.get(http::MINIO_INTERNAL_PREFIX.len()..)
|
|
.filter(|_| Self::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)
|
|
|| Self::starts_with_ignore_ascii_case(suffix, http::SUFFIX_REPLICATION_RESET_ARN_PREFIX)
|
|
// Raw compatibility keys are normalized and hashed separately below.
|
|
|| Self::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(super) 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_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");
|
|
assert_eq!(err, DiskError::ErasureReadQuorum);
|
|
}
|
|
|
|
#[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");
|
|
}
|
|
|
|
// 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 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"
|
|
);
|
|
}
|
|
}
|