fix(ecstore): keep degraded objects listable when drives are offline (#7010)

This commit is contained in:
唐小鸭
2026-09-01 20:16:57 +08:00
committed by GitHub
parent 394394cdfc
commit 43450df589
8 changed files with 483 additions and 41 deletions
+1 -1
View File
@@ -1 +1 @@
sha256=9c2b958035a038ffd5ab98cac5f59a1b8e6a16e141f109ec7fb956afc0f11105
sha256=51da41c54167602f2bd6c45921b39a44562bf3cfcdf468d992bb992c62cad7fd
+2 -2
View File
@@ -477,7 +477,7 @@ path = "junit.xml"
[profile.e2e-nightly]
default-filter = """
package(e2e_test)
& test(/^(admin_timeout_regression_test|cluster_concurrency_test|cluster_multidrive_pool_test|heal_erasure_disk_rebuild_test|namespace_lock_quorum_test|object_lambda_test|stale_multipart_cleanup_cluster_test)::/)
& test(/^(admin_timeout_regression_test|cluster_concurrency_test|cluster_multidrive_pool_test|degraded_listing_availability_test|heal_erasure_disk_rebuild_test|namespace_lock_quorum_test|object_lambda_test|stale_multipart_cleanup_cluster_test)::/)
"""
fail-fast = false
@@ -533,7 +533,7 @@ path = "junit.xml"
default-filter = """
package(e2e_test)
& !test(/^protocols::/)
& !test(/^(admin_timeout_regression_test|cluster_concurrency_test|cluster_multidrive_pool_test|heal_erasure_disk_rebuild_test|namespace_lock_quorum_test|object_lambda_test|stale_multipart_cleanup_cluster_test)::/)
& !test(/^(admin_timeout_regression_test|cluster_concurrency_test|cluster_multidrive_pool_test|degraded_listing_availability_test|heal_erasure_disk_rebuild_test|namespace_lock_quorum_test|object_lambda_test|stale_multipart_cleanup_cluster_test)::/)
& !test(/^replication_extension_test::/)
"""
fail-fast = false
@@ -0,0 +1,174 @@
// 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.
//! Regression: an object legally committed at degraded write quorum must stay
//! listable while a *different* drive is offline.
//!
//! On a 4-drive EC 2+2 set, a PUT made while one drive is down persists
//! `xl.meta` on 3 of 4 drives (write quorum). If a different drive later goes
//! offline before heal converges, a strict latest-listing quorum of 3 can only
//! ever observe 2 copies, so ListObjectsV2 silently dropped the object even
//! though GetObject (read quorum 2) still succeeded. Exposed by the flaky
//! "Mixed-version rolling upgrade from rc.2" CI lane (run 33478999853); the
//! product fix relaxes the listing's required object quorum by the number of
//! set drives the listing could not consult (see
//! `latest_listing_required_object_quorum` in
//! `crates/ecstore/src/store/list_objects.rs`).
#[cfg(test)]
mod tests {
use crate::common::{RustFSTestClusterEnvironment, init_logging};
use aws_sdk_s3::Client;
use bytes::Bytes;
use std::collections::HashSet;
use std::error::Error;
use std::time::{Duration, Instant};
use tracing::info;
type TestResult = Result<(), Box<dyn Error + Send + Sync>>;
const BUCKET: &str = "degraded-listing-availability";
const OBJECT_COUNT: usize = 8;
/// Well under the observed heal-convergence gap (~50s in the CI incident),
/// so a listing that only completes after heal restores the missing copy
/// still fails this deadline on a regressed build.
const LISTING_DEADLINE: Duration = Duration::from_secs(25);
const GET_RETRY_DEADLINE: Duration = Duration::from_secs(15);
const PUT_RETRY_DEADLINE: Duration = Duration::from_secs(15);
fn object_key(idx: usize) -> String {
format!("degraded-object-{idx:02}")
}
async fn list_all_keys(client: &Client) -> Result<HashSet<String>, Box<dyn Error + Send + Sync>> {
let mut keys = HashSet::new();
let mut continuation_token: Option<String> = None;
loop {
let response = client
.list_objects_v2()
.bucket(BUCKET)
.set_continuation_token(continuation_token.clone())
.send()
.await?;
keys.extend(
response
.contents()
.iter()
.filter_map(|object| object.key().map(str::to_owned)),
);
match response.next_continuation_token() {
Some(token) => continuation_token = Some(token.to_owned()),
None => break,
}
}
Ok(keys)
}
/// 4-node single-drive cluster (EC 2+2, write quorum 3):
/// 1. Stop node 1 and PUT objects — each commits on nodes {0, 2, 3} only.
/// 2. Stop node 3 (a holder drive), then bring node 1 back before heal can
/// recreate the missing copies there.
/// 3. Every object still satisfies read quorum (nodes 0 and 2), so GET
/// must succeed AND ListObjectsV2 must report every key well before
/// heal converges.
#[tokio::test]
async fn degraded_write_remains_listable_while_a_different_drive_is_offline() -> TestResult {
init_logging();
let mut cluster = RustFSTestClusterEnvironment::new(4).await?;
// Listing availability must not depend on heal convergence: disable
// the background healers so the degraded objects keep their metadata
// on exactly 3 of 4 drives for the whole test.
cluster.set_env("RUSTFS_HEAL_ENABLED", "false");
cluster.set_env("RUSTFS_SCANNER_ENABLED", "false");
cluster.start().await?;
cluster.create_test_bucket(BUCKET).await?;
let client = cluster.create_s3_client(0)?;
info!("stopping node 1 so the uploads commit at degraded write quorum (3 of 4)");
cluster.stop_node(1)?;
// The first writes after a node drops can see transient 503s while the
// survivors notice the dead peer; retry briefly (overwrites of the same
// unversioned key are idempotent).
for idx in 0..OBJECT_COUNT {
let key = object_key(idx);
let body = format!("degraded listing payload {idx}");
let deadline = Instant::now() + PUT_RETRY_DEADLINE;
loop {
let request = client
.put_object()
.bucket(BUCKET)
.key(&key)
.body(Bytes::from(body.clone()).into());
match request.send().await {
Ok(_) => break,
Err(error) if Instant::now() < deadline => {
info!("retrying degraded PUT for {key}: {error}");
tokio::time::sleep(Duration::from_millis(500)).await;
}
Err(error) => return Err(format!("degraded PUT for {key} failed: {error}").into()),
}
}
}
info!("stopping node 3 (holds a copy) and restoring node 1 (holds none)");
cluster.stop_node(3)?;
cluster.start_node(1).await?;
// The first requests after a node drops can see transient 503s while
// the survivors notice the dead peer; retry briefly before asserting.
for idx in 0..OBJECT_COUNT {
let key = object_key(idx);
let deadline = Instant::now() + GET_RETRY_DEADLINE;
let body = loop {
match client.get_object().bucket(BUCKET).key(&key).send().await {
Ok(response) => break response.body.collect().await?.into_bytes(),
Err(error) if Instant::now() < deadline => {
info!("retrying degraded GET for {key}: {error}");
tokio::time::sleep(Duration::from_millis(500)).await;
}
Err(error) => return Err(format!("degraded object {key} failed read quorum GET: {error}").into()),
}
};
assert!(!body.is_empty(), "degraded object {key} should read back at read quorum");
}
let expected: HashSet<String> = (0..OBJECT_COUNT).map(object_key).collect();
let deadline = Instant::now() + LISTING_DEADLINE;
let listed = loop {
let listed = match list_all_keys(&client).await {
Ok(keys) => keys,
Err(error) if Instant::now() < deadline => {
info!("retrying degraded listing: {error}");
tokio::time::sleep(Duration::from_millis(500)).await;
continue;
}
Err(error) => return Err(error),
};
if expected.is_subset(&listed) {
break listed;
}
assert!(
Instant::now() < deadline,
"objects readable at read quorum stayed missing from ListObjectsV2 for {LISTING_DEADLINE:?}: \
missing={:?} listed={listed:?}",
expected.difference(&listed).collect::<Vec<_>>(),
);
tokio::time::sleep(Duration::from_millis(500)).await;
};
info!(listed = listed.len(), "degraded objects are listable while node 3 is offline");
Ok(())
}
}
+5
View File
@@ -348,6 +348,11 @@ mod delete_regression_test;
#[cfg(test)]
mod listing_regression_test;
// Cluster regression: objects committed at degraded write quorum must stay
// listable while a different drive is offline (CI run 33478999853).
#[cfg(test)]
mod degraded_listing_availability_test;
// P1 regression: bucket statistics accuracy (rustfs#5615, #5008, #5116, #5055, #3898, #1012)
#[cfg(test)]
mod bucket_stats_regression_test;
+189 -33
View File
@@ -2522,6 +2522,7 @@ fn list_metadata_resolution_params(
listing_quorum: usize,
latest_object_quorum: usize,
versioned: bool,
write_quorum_slack: usize,
) -> MetadataResolutionParams {
let quorum = if versioned {
listing_quorum
@@ -2532,6 +2533,7 @@ fn list_metadata_resolution_params(
dir_quorum: quorum,
obj_quorum: quorum,
bucket,
write_quorum_slack,
..Default::default()
};
@@ -2832,7 +2834,10 @@ fn cached_entry_needs_supplement(
let mut selected_object_versions = 0;
for version in cached.versions.iter() {
let required_quorum = version.write_quorum(resolver.obj_quorum).max(resolver.obj_quorum);
let required_quorum = version
.write_quorum(resolver.obj_quorum)
.saturating_sub(resolver.write_quorum_slack)
.max(resolver.obj_quorum);
if version_requires_supplement(required_quorum, reader_disks, selected_object_versions, resolver.requested_versions) {
return true;
}
@@ -2983,7 +2988,10 @@ fn resolve_agreed_listing_entry(
let mut needs_supplement = false;
for (idx, version) in cached.versions.iter().enumerate() {
let required_quorum = version.write_quorum(resolver.obj_quorum).max(resolver.obj_quorum);
let required_quorum = version
.write_quorum(resolver.obj_quorum)
.saturating_sub(resolver.write_quorum_slack)
.max(resolver.obj_quorum);
if reader_disks < required_quorum {
needs_supplement |=
version_requires_supplement(required_quorum, reader_disks, selected_object_versions, resolver.requested_versions);
@@ -3037,21 +3045,51 @@ fn latest_listing_object_quorum(
drive_count: usize,
parity_count: usize,
enforce_write_quorum: bool,
unreachable_disks: usize,
) -> usize {
latest_listing_required_object_quorum(listing_quorum, drive_count, parity_count, enforce_write_quorum)
latest_listing_required_object_quorum(listing_quorum, drive_count, parity_count, enforce_write_quorum, unreachable_disks)
}
/// Object quorum a listed "latest" version must reach among the drives the
/// listing can actually consult.
///
/// An object legally committed at write quorum can have up to
/// `unreachable_disks` of its metadata copies on drives that are offline for
/// this listing, so the write-quorum requirement is relaxed by that amount.
/// The result is floored at the erasure read quorum (data drives): below that
/// the object could not be read back either, and a quorum-deleted object
/// leaves at most `drive_count - write_quorum < read_quorum` stale copies, so
/// the floor also keeps deleted objects from resurfacing.
///
/// Trade-off: while a drive is offline, a torn overwrite that reached only
/// `write_quorum - 1` drives becomes indistinguishable from a committed write
/// whose missing copy sits on the offline drive, so it can be listed as
/// latest. GET at read quorum serves that same version in that state, so the
/// listing stays consistent with reads instead of hiding readable objects.
fn latest_listing_required_object_quorum(
listing_quorum: usize,
drive_count: usize,
parity_count: usize,
enforce_write_quorum: bool,
unreachable_disks: usize,
) -> usize {
if !enforce_write_quorum {
return listing_quorum;
}
write_quorum_for_drive_count(drive_count, parity_count).max(listing_quorum)
let read_quorum = drive_count.saturating_sub(parity_count);
write_quorum_for_drive_count(drive_count, parity_count)
.saturating_sub(unreachable_disks)
.max(read_quorum)
.max(listing_quorum)
}
fn latest_listing_write_quorum_slack(enforce_write_quorum: bool, drive_count: usize, online_disks: usize) -> usize {
if !enforce_write_quorum {
return 0;
}
drive_count.saturating_sub(online_disks)
}
fn enforce_latest_listing_write_quorum(strict_latest: bool, ask_disks: &str) -> bool {
@@ -4371,6 +4409,9 @@ impl ECStore {
for eset in self.pools.iter() {
for set in eset.disk_set.iter() {
let (mut disks, infos, _) = set.get_online_disks_with_healing_and_info(true).await;
// Captured before any quorum-based filtering: only genuinely
// unreachable drives may relax the write-quorum requirement.
let online_disks = disks.len();
let opts = opts.clone();
let (sender, list_out_rx) = mpsc::channel::<MetaCacheEntry>(1);
@@ -4396,11 +4437,14 @@ impl ECStore {
let listing_quorum = listing_quorum_from_ask_disks(ask_disks);
let enforce_write_quorum = enforce_latest_listing_write_quorum(opts.latest_only, &opts.ask_disks);
let write_quorum_parity = set.default_parity_count;
let write_quorum_slack =
latest_listing_write_quorum_slack(enforce_write_quorum, set.set_drive_count, online_disks);
let required_obj_quorum = latest_listing_required_object_quorum(
listing_quorum,
set.set_drive_count,
write_quorum_parity,
enforce_write_quorum,
write_quorum_slack,
);
ask_disks = expand_ask_disks_for_object_quorum(ask_disks, disks.len(), required_obj_quorum);
let fallback_disks = {
@@ -4422,11 +4466,17 @@ impl ECStore {
set.set_drive_count,
write_quorum_parity,
enforce_write_quorum,
write_quorum_slack,
);
let raw_min_disks = latest_listing_raw_min_disks(listing_quorum, obj_quorum, enforce_write_quorum);
let resolver =
list_metadata_resolution_params(bucket.to_owned(), listing_quorum, obj_quorum, !opts.latest_only);
let resolver = list_metadata_resolution_params(
bucket.to_owned(),
listing_quorum,
obj_quorum,
!opts.latest_only,
write_quorum_slack,
);
let agreed_resolver = resolver.clone();
let partial_resolver = resolver.clone();
let reader_disks = disks.len();
@@ -5613,6 +5663,9 @@ impl Sets {
for set in &self.disk_set {
let (mut disks, infos, _) = set.get_online_disks_with_healing_and_info(true).await;
// Captured before any quorum-based filtering: only genuinely
// unreachable drives may relax the write-quorum requirement.
let online_disks = disks.len();
let opts = opts.clone();
let (sender, list_out_rx) = mpsc::channel::<MetaCacheEntry>(1);
inputs.push(list_out_rx);
@@ -5638,11 +5691,14 @@ impl Sets {
let listing_quorum = listing_quorum_from_ask_disks(ask_disks);
let enforce_write_quorum = enforce_latest_listing_write_quorum(opts.latest_only, &opts.ask_disks);
let write_quorum_parity = set.default_parity_count;
let write_quorum_slack =
latest_listing_write_quorum_slack(enforce_write_quorum, set.set_drive_count, online_disks);
let required_obj_quorum = latest_listing_required_object_quorum(
listing_quorum,
set.set_drive_count,
write_quorum_parity,
enforce_write_quorum,
write_quorum_slack,
);
ask_disks = expand_ask_disks_for_object_quorum(ask_disks, disks.len(), required_obj_quorum);
let fallback_disks = if let Some(asked_disks) = positive_ask_disks(ask_disks)
@@ -5657,10 +5713,21 @@ impl Sets {
let fallback_disks = Arc::new(fallback_disks);
let claim_tracker = FallbackClaimTracker::default();
let obj_quorum =
latest_listing_object_quorum(listing_quorum, set.set_drive_count, write_quorum_parity, enforce_write_quorum);
let obj_quorum = latest_listing_object_quorum(
listing_quorum,
set.set_drive_count,
write_quorum_parity,
enforce_write_quorum,
write_quorum_slack,
);
let raw_min_disks = latest_listing_raw_min_disks(listing_quorum, obj_quorum, enforce_write_quorum);
let resolver = list_metadata_resolution_params(bucket.to_owned(), listing_quorum, obj_quorum, !opts.latest_only);
let resolver = list_metadata_resolution_params(
bucket.to_owned(),
listing_quorum,
obj_quorum,
!opts.latest_only,
write_quorum_slack,
);
let agreed_resolver = resolver.clone();
let partial_resolver = resolver.clone();
let reader_disks = disks.len();
@@ -6582,6 +6649,9 @@ impl SetDisks {
let list_path_started = std::time::Instant::now();
let (mut disks, infos, _) = self.get_online_disks_with_healing_and_info(true).await;
// Captured before any quorum-based filtering: only genuinely
// unreachable drives may relax the write-quorum requirement.
let online_disks = disks.len();
let mut ask_disks = get_list_quorum(&opts.ask_disks, self.set_drive_count as i32);
if ask_disks == -1 {
@@ -6605,11 +6675,13 @@ impl SetDisks {
let enforce_write_quorum = enforce_latest_listing_write_quorum(!opts.versioned, &opts.ask_disks);
let write_quorum_parity = self.default_parity_count;
let write_quorum_slack = latest_listing_write_quorum_slack(enforce_write_quorum, self.set_drive_count, online_disks);
let required_obj_quorum = latest_listing_required_object_quorum(
listing_quorum,
self.set_drive_count,
write_quorum_parity,
enforce_write_quorum,
write_quorum_slack,
);
ask_disks = expand_ask_disks_for_object_quorum(ask_disks, disks.len(), required_obj_quorum);
let mut fallback_disks = Vec::new();
@@ -6627,10 +6699,21 @@ impl SetDisks {
let bucket = opts.bucket.clone();
let base_dir = opts.base_dir.clone();
let latest_object_quorum =
latest_listing_object_quorum(listing_quorum, self.set_drive_count, write_quorum_parity, enforce_write_quorum);
let latest_object_quorum = latest_listing_object_quorum(
listing_quorum,
self.set_drive_count,
write_quorum_parity,
enforce_write_quorum,
write_quorum_slack,
);
let raw_min_disks = latest_listing_raw_min_disks(listing_quorum, latest_object_quorum, enforce_write_quorum);
let resolver = list_metadata_resolution_params(bucket.clone(), listing_quorum, latest_object_quorum, opts.versioned);
let resolver = list_metadata_resolution_params(
bucket.clone(),
listing_quorum,
latest_object_quorum,
opts.versioned,
write_quorum_slack,
);
let agreed_resolver = resolver.clone();
let partial_resolver = resolver.clone();
let reader_disks = disks.len();
@@ -6642,6 +6725,7 @@ impl SetDisks {
asked_disks = ask_disks,
listing_quorum = listing_quorum,
latest_object_quorum = latest_object_quorum,
write_quorum_slack = write_quorum_slack,
raw_min_disks = raw_min_disks,
fallback_disks = fallback_disks.len(),
limit = opts.limit,
@@ -6889,10 +6973,11 @@ mod test {
PERSISTENT_KEY_ONLY_INDEX_CHECKPOINT_HEADER, PERSISTENT_KEY_ONLY_INDEX_FORMAT_VERSION,
PERSISTENT_KEY_ONLY_INDEX_GENERATION_HEADER, PERSISTENT_KEY_ONLY_INDEX_HEADER, PersistentKeyOnlyIndex,
PersistentListMetadataObject, RUSTFS_META_BUCKET, VerifiedIndexCandidateStats, VersionMarker,
current_list_objects_mutation_sequence, encode_persistent_list_metadata_object, enforce_latest_listing_write_quorum,
expand_ask_disks_for_object_quorum, fallback_entries_for_object, gather_results, latest_listing_allow_agreed_objects,
latest_listing_object_quorum, latest_listing_raw_min_disks, latest_listing_required_object_quorum, list_marker_key,
list_merged_entry_channel, list_metadata_resolution_params, list_objects_from_metadata_snapshot_candidates,
cached_entry_needs_supplement, current_list_objects_mutation_sequence, encode_persistent_list_metadata_object,
enforce_latest_listing_write_quorum, expand_ask_disks_for_object_quorum, fallback_entries_for_object, gather_results,
latest_listing_allow_agreed_objects, latest_listing_object_quorum, latest_listing_raw_min_disks,
latest_listing_required_object_quorum, latest_listing_write_quorum_slack, list_marker_key, list_merged_entry_channel,
list_metadata_resolution_params, list_objects_from_metadata_snapshot_candidates,
list_objects_from_verified_index_candidates, list_objects_from_verified_index_candidates_with_optional_stats,
list_objects_from_verified_index_candidates_with_stats, list_objects_index_mode_from_env,
list_objects_index_provider_from_env, list_objects_index_provider_state_from_env,
@@ -9113,7 +9198,7 @@ mod test {
#[test]
fn list_metadata_resolution_params_limits_plain_listing_to_latest_version() {
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, 3, false);
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, 3, false, 0);
assert_eq!(resolver.dir_quorum, 3);
assert_eq!(resolver.obj_quorum, 3);
@@ -9123,7 +9208,7 @@ mod test {
#[test]
fn list_metadata_resolution_params_keeps_all_versions_for_version_listing() {
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, true);
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, true, 0);
assert_eq!(resolver.dir_quorum, 3);
assert_eq!(resolver.obj_quorum, 3);
@@ -9133,22 +9218,22 @@ mod test {
#[test]
fn latest_listing_object_quorum_uses_write_quorum_for_strict_latest_listing() {
let required_quorum = latest_listing_required_object_quorum(2, 8, 4, true);
let required_quorum = latest_listing_required_object_quorum(2, 8, 4, true, 0);
let ask_disks = expand_ask_disks_for_object_quorum(4, 8, required_quorum);
assert_eq!(required_quorum, 5);
assert_eq!(ask_disks, 5);
assert_eq!(latest_listing_object_quorum(2, 8, 4, true), 5);
assert_eq!(latest_listing_object_quorum(2, 8, 4, true, 0), 5);
}
#[test]
fn latest_listing_object_quorum_calculates_low_parity_write_quorum() {
let required_quorum = latest_listing_required_object_quorum(2, 8, 1, true);
let required_quorum = latest_listing_required_object_quorum(2, 8, 1, true, 0);
let ask_disks = expand_ask_disks_for_object_quorum(4, 8, required_quorum);
assert_eq!(required_quorum, 7);
assert_eq!(ask_disks, 7);
assert_eq!(latest_listing_object_quorum(2, 8, 1, true), 7);
assert_eq!(latest_listing_object_quorum(2, 8, 1, true, 0), 7);
}
#[test]
@@ -9157,21 +9242,92 @@ mod test {
assert!(!enforce_latest_listing_write_quorum(true, "disk"));
assert!(enforce_latest_listing_write_quorum(true, "optimal"));
assert!(!enforce_latest_listing_write_quorum(false, "optimal"));
assert_eq!(latest_listing_required_object_quorum(1, 4, 2, false), 1);
assert_eq!(latest_listing_object_quorum(1, 4, 2, false), 1);
assert_eq!(latest_listing_required_object_quorum(2, 4, 2, false), 2);
assert_eq!(latest_listing_object_quorum(2, 4, 2, false), 2);
assert_eq!(latest_listing_required_object_quorum(1, 4, 2, false, 0), 1);
assert_eq!(latest_listing_object_quorum(1, 4, 2, false, 0), 1);
assert_eq!(latest_listing_required_object_quorum(2, 4, 2, false, 0), 2);
assert_eq!(latest_listing_object_quorum(2, 4, 2, false, 0), 2);
assert_eq!(expand_ask_disks_for_object_quorum(2, 4, 2), 2);
}
#[test]
fn latest_listing_object_quorum_relaxes_write_quorum_by_unreachable_drives() {
// 4-drive set, EC 2+2: with one drive offline an object committed at
// write quorum 3 can only ever show 2 metadata copies to the listing.
let slack = latest_listing_write_quorum_slack(true, 4, 3);
assert_eq!(slack, 1);
assert_eq!(latest_listing_required_object_quorum(2, 4, 2, true, slack), 2);
assert_eq!(latest_listing_required_object_quorum(2, 4, 2, true, 0), 3);
// The relaxed quorum never drops below the erasure read quorum, so a
// quorum-deleted object (at most one stale copy) stays hidden even
// with half the set unreachable.
let slack = latest_listing_write_quorum_slack(true, 4, 2);
assert_eq!(slack, 2);
assert_eq!(latest_listing_required_object_quorum(1, 4, 2, true, slack), 2);
assert_eq!(latest_listing_write_quorum_slack(false, 4, 2), 0);
}
#[test]
fn latest_listing_resolves_degraded_object_when_a_set_drive_is_unreachable() {
let mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let entry = test_object_meta_entry_with_erasure_versions("object", &[(mod_time, "etag", 2, 2)]);
let entries = || MetaCacheEntries(vec![Some(entry.clone()), Some(entry.clone()), None]);
let slack = latest_listing_write_quorum_slack(true, 4, 3);
let obj_quorum = latest_listing_object_quorum(2, 4, 2, true, slack);
assert_eq!(obj_quorum, 2);
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, obj_quorum, false, slack);
let resolved = resolve_listing_entries(entries(), resolver, true)
.expect("object committed at write quorum should stay listable with one holder drive offline");
assert_eq!(resolved.name, "object");
// Without the unreachable-drive slack the same sample is dropped even
// though the object still satisfies read quorum for GET.
let strict_quorum = latest_listing_object_quorum(2, 4, 2, true, 0);
let strict_resolver = list_metadata_resolution_params("bucket".to_string(), 2, strict_quorum, false, 0);
assert!(resolve_listing_entries(entries(), strict_resolver, true).is_none());
}
#[test]
fn agreed_listing_entry_relaxes_write_quorum_by_unreachable_drives() {
let mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let entry = test_object_meta_entry_with_erasure_versions("object", &[(mod_time, "etag", 2, 2)]);
let mut resolver = list_metadata_resolution_params("bucket".to_string(), 2, 2, false, 1);
assert!(matches!(
resolve_agreed_listing_entry(entry.clone(), 2, resolver.clone(), true),
ListingEntryResolution::Resolved(_)
));
resolver.write_quorum_slack = 0;
assert!(matches!(
resolve_agreed_listing_entry(entry, 2, resolver, true),
ListingEntryResolution::NeedsSupplement(_, _)
));
}
#[test]
fn cached_entry_supplement_check_honors_write_quorum_slack() {
let mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let mut entry = test_object_meta_entry_with_erasure_versions("object", &[(mod_time, "etag", 2, 2)]);
let cached = entry.xl_meta().expect("test entry should decode");
let mut resolver = list_metadata_resolution_params("bucket".to_string(), 2, 2, false, 1);
assert!(!cached_entry_needs_supplement(&cached, 2, &resolver, true));
resolver.write_quorum_slack = 0;
assert!(cached_entry_needs_supplement(&cached, 2, &resolver, true));
}
#[test]
fn latest_listing_object_quorum_requires_write_quorum_when_degraded_cannot_satisfy_it() {
let required_quorum = latest_listing_required_object_quorum(2, 8, 4, true);
let required_quorum = latest_listing_required_object_quorum(2, 8, 4, true, 0);
let ask_disks = expand_ask_disks_for_object_quorum(4, 4, required_quorum);
assert_eq!(required_quorum, 5);
assert_eq!(ask_disks, 4);
assert_eq!(latest_listing_object_quorum(2, 8, 4, true), 5);
assert_eq!(latest_listing_object_quorum(2, 8, 4, true, 0), 5);
}
#[tokio::test]
@@ -9186,7 +9342,7 @@ mod test {
Some(deleted.clone()),
Some(deleted.clone()),
]);
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, 4, false);
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, 4, false, 0);
assert_eq!(listing_entries_supplement_target(&entries, &resolver, true).as_deref(), Some("object"));
assert_eq!(listing_entries_supplement_target(&entries, &resolver, false), None);
@@ -9251,7 +9407,7 @@ mod test {
let delete_mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_400).expect("valid timestamp");
let stale = test_object_meta_entry_with_erasure_versions("object", &[(object_mod_time, "object-etag", 4, 2)]);
let deleted = test_object_with_delete_marker_meta_entry("object", object_mod_time, delete_mod_time);
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, 4, false);
let resolver = list_metadata_resolution_params("bucket".to_string(), 2, 4, false, 0);
let entries = MetaCacheEntries(vec![
Some(stale.clone()),
Some(stale.clone()),
@@ -9352,7 +9508,7 @@ mod test {
&[(old_mod_time, "old-etag", 4, 4), (new_mod_time, "new-etag", 7, 1)],
);
let fallback_old_entry = test_object_meta_entry_with_erasure_versions("object", &[(old_mod_time, "old-etag", 4, 4)]);
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, false);
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, false, 0);
let seen = Arc::new(Mutex::new(Vec::new()));
let seen_clone = seen.clone();
@@ -9407,7 +9563,7 @@ mod test {
"object",
&[(old_mod_time, "old-etag", 4, 4), (new_mod_time, "new-etag", 7, 1)],
);
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, false);
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, false, 0);
let seen = Arc::new(Mutex::new(Vec::new()));
let seen_clone = seen.clone();
@@ -9451,7 +9607,7 @@ mod test {
let new_mod_time = time::OffsetDateTime::from_unix_timestamp(1_705_312_400).expect("valid timestamp");
let entry = test_object_meta_entry_with_erasure_versions("object", &[(new_mod_time, "new-etag", 7, 1)]);
let fallback_entry = entry.clone();
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, false);
let resolver = list_metadata_resolution_params("bucket".to_string(), 3, 5, false, 0);
let seen = Arc::new(Mutex::new(Vec::new()));
let seen_clone = seen.clone();
+6 -4
View File
@@ -3235,16 +3235,17 @@ pub fn merge_file_meta_versions(
requested_versions: usize,
versions: &[Vec<FileMetaShallowVersion>],
) -> Vec<FileMetaShallowVersion> {
merge_file_meta_versions_inner(quorum, strict, requested_versions, false, versions)
merge_file_meta_versions_inner(quorum, strict, requested_versions, false, 0, versions)
}
pub(crate) fn merge_file_meta_versions_with_write_quorum(
quorum: usize,
strict: bool,
requested_versions: usize,
write_quorum_slack: usize,
versions: &[Vec<FileMetaShallowVersion>],
) -> Vec<FileMetaShallowVersion> {
merge_file_meta_versions_inner(quorum, strict, requested_versions, true, versions)
merge_file_meta_versions_inner(quorum, strict, requested_versions, true, write_quorum_slack, versions)
}
fn merge_file_meta_versions_inner(
@@ -3252,6 +3253,7 @@ fn merge_file_meta_versions_inner(
mut strict: bool,
requested_versions: usize,
enforce_write_quorum: bool,
write_quorum_slack: usize,
versions: &[Vec<FileMetaShallowVersion>],
) -> Vec<FileMetaShallowVersion> {
if quorum == 0 {
@@ -3269,7 +3271,7 @@ fn merge_file_meta_versions_inner(
let required_quorum = versions[0]
.first()
.map(|version| version.write_quorum(quorum).max(quorum))
.map(|version| version.write_quorum(quorum).saturating_sub(write_quorum_slack).max(quorum))
.unwrap_or(quorum);
if versions.len() >= required_quorum {
return versions[0].clone();
@@ -3283,7 +3285,7 @@ fn merge_file_meta_versions_inner(
let required_quorum = |version: &FileMetaShallowVersion| {
if enforce_write_quorum {
version.write_quorum(quorum).max(quorum)
version.write_quorum(quorum).saturating_sub(write_quorum_slack).max(quorum)
} else {
quorum
}
+105 -1
View File
@@ -53,6 +53,12 @@ pub struct MetadataResolutionParams {
pub requested_versions: usize,
pub bucket: String,
pub strict: bool,
/// Number of set drives that were unreachable when the listing snapshot was
/// taken. Write-quorum enforcement relaxes each version's required quorum by
/// this amount (never below `obj_quorum`): a version legally committed at
/// write quorum can have that many of its metadata copies on drives no
/// reader could consult, and must not be dropped for it.
pub write_quorum_slack: usize,
pub candidates: Vec<Vec<FileMetaShallowVersion>>,
}
@@ -416,6 +422,7 @@ impl MetaCacheEntries {
requested_versions: 0,
bucket: bucket.to_string(),
strict: false,
write_quorum_slack: 0,
candidates: Vec::new(),
})
}
@@ -693,7 +700,12 @@ impl MetaCacheEntries {
.cached
.as_ref()
.and_then(|cached| cached.versions.first())
.map(|version| version.write_quorum(params.obj_quorum).max(params.obj_quorum))
.map(|version| {
version
.write_quorum(params.obj_quorum)
.saturating_sub(params.write_quorum_slack)
.max(params.obj_quorum)
})
.unwrap_or(params.obj_quorum)
} else {
params.obj_quorum
@@ -720,6 +732,7 @@ impl MetaCacheEntries {
params.obj_quorum,
params.strict,
params.requested_versions,
params.write_quorum_slack,
&params.candidates,
)
} else {
@@ -2415,6 +2428,97 @@ mod tests {
assert_eq!(info.metadata.get("etag").map(String::as_str), Some("old-etag"));
}
#[test]
fn resolve_with_write_quorum_relaxes_requirement_by_unreachable_drives() {
let mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_400).expect("valid timestamp");
// EC 2+2 version, write quorum 3 of 4: committed with one holder drive
// now offline, so only 2 of the 3 reachable drives carry the metadata.
let entry = metacache_entry_with_erasure(mod_time, "etag", 2, 2);
let entries = || MetaCacheEntries(vec![Some(entry.clone()), Some(entry.clone()), None]);
let resolved = entries()
.resolve_with_write_quorum(MetadataResolutionParams {
obj_quorum: 2,
requested_versions: 1,
bucket: "bucket".to_string(),
strict: true,
write_quorum_slack: 1,
..Default::default()
})
.expect("committed version should resolve when the missing copy sits on an unreachable drive");
let info = resolved
.to_fileinfo("bucket")
.expect("resolved committed metadata should decode as file info");
assert_eq!(info.mod_time, Some(mod_time));
// Without the slack the same sample is rejected outright.
let rejected = entries().resolve_with_write_quorum(MetadataResolutionParams {
obj_quorum: 3,
requested_versions: 1,
bucket: "bucket".to_string(),
strict: true,
..Default::default()
});
assert!(rejected.is_none());
}
#[test]
fn resolve_with_write_quorum_slack_accepts_committed_latest_during_merge() {
let old_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let new_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_400).expect("valid timestamp");
// The newer EC 2+2 version committed at write quorum 3 of 4; one holder
// drive is offline, and a third reachable drive still carries only the
// previous version. The candidates disagree, so the merge path (not the
// all-agree branch) must honor the slack.
let old_entry = metacache_entry_with_erasure(old_mod_time, "old-etag", 2, 2);
let new_and_old_entry =
metacache_entry_with_erasure_versions(&[(old_mod_time, "old-etag", 2, 2), (new_mod_time, "new-etag", 2, 2)]);
let resolved = MetaCacheEntries(vec![Some(new_and_old_entry.clone()), Some(new_and_old_entry), Some(old_entry)])
.resolve_with_write_quorum(MetadataResolutionParams {
obj_quorum: 2,
requested_versions: 1,
bucket: "bucket".to_string(),
strict: true,
write_quorum_slack: 1,
..Default::default()
})
.expect("committed latest should survive the merge when its missing copy is on an unreachable drive");
let info = resolved
.to_fileinfo("bucket")
.expect("resolved committed metadata should decode as file info");
assert_eq!(info.mod_time, Some(new_mod_time));
assert_eq!(info.metadata.get("etag").map(String::as_str), Some("new-etag"));
}
#[test]
fn resolve_with_write_quorum_slack_keeps_partial_latest_hidden_during_merge() {
let old_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_300).expect("valid timestamp");
let new_mod_time = OffsetDateTime::from_unix_timestamp(1_705_312_400).expect("valid timestamp");
// The newer EC 2+2 version (write quorum 3) reached only ONE drive: even
// with one drive unreachable (slack 1) it cannot prove the relaxed
// quorum of 2, so the committed previous version must win the merge.
let old_entry = metacache_entry_with_erasure(old_mod_time, "old-etag", 2, 2);
let new_and_old_entry =
metacache_entry_with_erasure_versions(&[(old_mod_time, "old-etag", 2, 2), (new_mod_time, "new-etag", 2, 2)]);
let resolved = MetaCacheEntries(vec![Some(new_and_old_entry), Some(old_entry.clone()), Some(old_entry)])
.resolve_with_write_quorum(MetadataResolutionParams {
obj_quorum: 2,
requested_versions: 1,
bucket: "bucket".to_string(),
strict: true,
write_quorum_slack: 1,
..Default::default()
})
.expect("committed previous version should resolve after rejecting the partial latest");
let info = resolved
.to_fileinfo("bucket")
.expect("resolved committed metadata should decode as file info");
assert_eq!(info.mod_time, Some(old_mod_time));
assert_eq!(info.metadata.get("etag").map(String::as_str), Some("old-etag"));
}
#[test]
fn resolve_rejects_partial_directory_below_dir_quorum() {
let partial_dir = metacache_dir_entry("prefix/");
+1
View File
@@ -43,6 +43,7 @@
| copy_source_invalid_date_test | 1 | ✅ |
| create_bucket_region_test | 2 | ✅ |
| data_usage_test | 2 | |
| degraded_listing_availability_test | 1 | 🌙 |
| degraded_read_eof_regression_test | 3 | |
| delete_marker_migration_semantics_test | 2 | ✅ |
| delete_object_no_content_length_test | 1 | |