diff --git a/.config/e2e-nightly-selection.txt b/.config/e2e-nightly-selection.txt index a3a2c14ad..6ed78d7c1 100644 --- a/.config/e2e-nightly-selection.txt +++ b/.config/e2e-nightly-selection.txt @@ -1 +1 @@ -sha256=9c2b958035a038ffd5ab98cac5f59a1b8e6a16e141f109ec7fb956afc0f11105 +sha256=51da41c54167602f2bd6c45921b39a44562bf3cfcdf468d992bb992c62cad7fd diff --git a/.config/nextest.toml b/.config/nextest.toml index b8ed8e626..d4da98b79 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -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 diff --git a/crates/e2e_test/src/degraded_listing_availability_test.rs b/crates/e2e_test/src/degraded_listing_availability_test.rs new file mode 100644 index 000000000..53929ce4d --- /dev/null +++ b/crates/e2e_test/src/degraded_listing_availability_test.rs @@ -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>; + + 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, Box> { + let mut keys = HashSet::new(); + let mut continuation_token: Option = 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 = (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::>(), + ); + tokio::time::sleep(Duration::from_millis(500)).await; + }; + info!(listed = listed.len(), "degraded objects are listable while node 3 is offline"); + + Ok(()) + } +} diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index c69e9ffea..e4720c060 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -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; diff --git a/crates/ecstore/src/store/list_objects.rs b/crates/ecstore/src/store/list_objects.rs index 381cd1cee..a426e9e64 100644 --- a/crates/ecstore/src/store/list_objects.rs +++ b/crates/ecstore/src/store/list_objects.rs @@ -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::(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::(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(); diff --git a/crates/filemeta/src/filemeta/version.rs b/crates/filemeta/src/filemeta/version.rs index 034e76017..3a8d69da9 100644 --- a/crates/filemeta/src/filemeta/version.rs +++ b/crates/filemeta/src/filemeta/version.rs @@ -3235,16 +3235,17 @@ pub fn merge_file_meta_versions( requested_versions: usize, versions: &[Vec], ) -> Vec { - 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], ) -> Vec { - 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], ) -> Vec { 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 } diff --git a/crates/filemeta/src/metacache.rs b/crates/filemeta/src/metacache.rs index 0571bc466..2c02a7bee 100644 --- a/crates/filemeta/src/metacache.rs +++ b/crates/filemeta/src/metacache.rs @@ -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>, } @@ -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, ¶ms.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/"); diff --git a/docs/testing/e2e-suite-inventory.md b/docs/testing/e2e-suite-inventory.md index 8a52be5af..cd9ddc69f 100644 --- a/docs/testing/e2e-suite-inventory.md +++ b/docs/testing/e2e-suite-inventory.md @@ -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 | |