fix(ecstore): preserve parity reserves for data-only GET (#6888)

fix(ecstore): hedge data-only GET with parity

Route the opt-in data-shards-only lockstep path through the bounded parity race and preserve deferred parity reserves across canceled hedges.

Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
houseme
2026-08-30 20:16:32 +08:00
committed by GitHub
parent 51532e19fb
commit 3d24526704
3 changed files with 103 additions and 65 deletions
+75 -45
View File
@@ -574,6 +574,7 @@ pub(crate) struct ParallelReader<R> {
read_timeout: Duration, read_timeout: Duration,
verify_reconstruction: bool, verify_reconstruction: bool,
locality_preference_enabled: bool, locality_preference_enabled: bool,
demand_bound_lockstep: bool,
// Request-scoped shard buffers keyed by shard index. Keeping ownership in // Request-scoped shard buffers keyed by shard index. Keeping ownership in
// `ParallelReader` avoids dropping unused parity/backup slot buffers between stripes. // `ParallelReader` avoids dropping unused parity/backup slot buffers between stripes.
buffers: ShardBufferPool, buffers: ShardBufferPool,
@@ -585,10 +586,8 @@ pub(crate) struct ParallelReader<R> {
// it to the current stripe when it is engaged mid-object (backlog#923). // it to the current stripe when it is engaged mid-object (backlog#923).
engaged: SmallVec<[bool; INLINE_SHARD_SLOTS]>, engaged: SmallVec<[bool; INLINE_SHARD_SLOTS]>,
deferred_handles: Vec<Option<DeferredReaderStripeHandle>>, deferred_handles: Vec<Option<DeferredReaderStripeHandle>>,
// Copy-source hedges use a fresh deferred reader so cancelling a hedge // Demand-bound hedges use a fresh deferred reader so cancelling a hedge
// never consumes the unopened reader reserved for a later stripe. The // never consumes the unopened reader reserved for a later stripe.
// vector is empty for callers that do not provide a reopen factory (tests
// and the ordinary GET path retain the handle-based behavior).
deferred_reopeners: Vec<Option<DeferredReaderReopener<R>>>, deferred_reopeners: Vec<Option<DeferredReaderReopener<R>>>,
stripe_index: usize, stripe_index: usize,
} }
@@ -777,9 +776,9 @@ where
// reads all live readers on every stripe — the pre-backlog#923 // reads all live readers on every stripe — the pre-backlog#923
// behavior. With the gate on, only data slots start engaged; parity is // behavior. With the gate on, only data slots start engaged; parity is
// engaged on demand, stripe-aligned through its deferred handle. // engaged on demand, stripe-aligned through its deferred handle.
let data_shards_only = get_lockstep_data_shards_only_enabled(); let demand_bound_lockstep = get_lockstep_data_shards_only_enabled();
let engaged: SmallVec<_> = (0..readers.len()) let engaged: SmallVec<_> = (0..readers.len())
.map(|index| !data_shards_only || index < e.data_shards) .map(|index| !demand_bound_lockstep || index < e.data_shards)
.collect(); .collect();
ParallelReader { ParallelReader {
readers, readers,
@@ -793,6 +792,7 @@ where
read_timeout, read_timeout,
verify_reconstruction, verify_reconstruction,
locality_preference_enabled: get_shard_locality_preference_enabled(), locality_preference_enabled: get_shard_locality_preference_enabled(),
demand_bound_lockstep,
buffers: ShardBufferPool::new(e.data_shards + e.parity_shards), buffers: ShardBufferPool::new(e.data_shards + e.parity_shards),
stripe_state: None, stripe_state: None,
engaged, engaged,
@@ -1275,7 +1275,7 @@ where
/// realigned (no pending deferred handle) is likewise retired instead of /// realigned (no pending deferred handle) is likewise retired instead of
/// being read out of position. /// being read out of position.
async fn read_lockstep(&mut self, state: &mut StripeReadState) { async fn read_lockstep(&mut self, state: &mut StripeReadState) {
if matches!(decode_read_policy(), DecodeReadPolicy::DemandBound) { if self.demand_bound_lockstep {
self.read_lockstep_demand_bound(state).await; self.read_lockstep_demand_bound(state).await;
return; return;
} }
@@ -1531,17 +1531,18 @@ where
} }
} }
/// Demand-bound lockstep stripe read used by server-side copy sources. /// Demand-bound data-shards-only lockstep stripe read.
/// ///
/// The ordinary lockstep path can cancel every in-flight reader once it /// The ordinary lockstep path can cancel every in-flight reader once it
/// has a quorum because all of its parity readers are already engaged. /// has a quorum because all of its parity readers are already engaged.
/// Copy sources keep parity unopened until a data reader is missing. A /// Copy sources and the data-shards-only rollout gate keep parity unopened
/// hedge therefore has to race the deferred parity reads against the /// until a data reader is missing. A hedge therefore has to race the
/// original data reads and may retire the latter only after the parity has /// deferred parity reads against the original data reads and may retire the
/// produced an actual decode-plus-verification quorum. The futures own /// latter only after parity has produced an actual decode-plus-verification
/// their readers so disjoint data/parity slots can be admitted while the /// quorum. The futures own their readers so disjoint data/parity slots can
/// other group is still pending; dropping an abandoned future retires its /// be admitted while the other group is still pending; dropping an
/// stream without leaving a borrowed slot behind. /// abandoned future retires its stream without leaving a borrowed slot
/// behind.
async fn read_lockstep_demand_bound(&mut self, state: &mut StripeReadState) { async fn read_lockstep_demand_bound(&mut self, state: &mut StripeReadState) {
let num_readers = self.readers.len(); let num_readers = self.readers.len();
state.reset(num_readers, self.data_shards); state.reset(num_readers, self.data_shards);
@@ -1576,14 +1577,14 @@ where
let mut completed = 0usize; let mut completed = 0usize;
let mut failed = 0usize; let mut failed = 0usize;
let mut first_shard_recorded = false; let mut first_shard_recorded = false;
let mut active = vec![false; num_readers]; let mut active: ActiveReaders = smallvec![false; num_readers];
let mut temporary_parity = vec![false; num_readers]; let mut temporary_parity: ActiveReaders = smallvec![false; num_readers];
// A deferred parity slot is attempted at most once per stripe. A // A deferred parity slot is attempted at most once per stripe. A
// failed disposable hedge keeps its unopened reserve for the next // failed disposable hedge keeps its unopened reserve for the next
// stripe, but must not be relaunched in a tight same-stripe retry // stripe, but must not be relaunched in a tight same-stripe retry
// loop (which would defeat the bounded fan-out and amplify a remote // loop (which would defeat the bounded fan-out and amplify a remote
// outage). // outage).
let mut attempted_parity = vec![false; num_readers]; let mut attempted_parity: ActiveReaders = smallvec![false; num_readers];
// Once a data reader has returned an error (or was already missing at // Once a data reader has returned an error (or was already missing at
// setup), the loss is permanent for lockstep alignment. Use the // setup), the loss is permanent for lockstep alignment. Use the
// deferred handle and keep parity engaged across subsequent stripes; // deferred handle and keep parity engaged across subsequent stripes;
@@ -4911,6 +4912,24 @@ mod tests {
/// read timeout even though both parity readers were available to engage. /// read timeout even though both parity readers were available to engage.
#[tokio::test] #[tokio::test]
async fn test_demand_bound_lockstep_hedges_to_deferred_parity_quorum() { async fn test_demand_bound_lockstep_hedges_to_deferred_parity_quorum() {
with_decode_read_policy(DecodeReadPolicy::DemandBound, assert_deferred_parity_hedges_slow_data()).await;
}
/// The ordinary GET rollout gate must use the same bounded parity race as
/// CopySource. Leaving it on the legacy lockstep loop deadlocks the hedge:
/// that loop waits for a parity success before cancelling the slow data
/// read, but does not admit deferred parity until after the data read ends.
#[tokio::test]
#[serial_test::serial]
async fn test_data_shards_only_gate_hedges_to_deferred_parity_quorum() {
temp_env::async_with_vars(
[(ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE, Some("true"))],
assert_deferred_parity_hedges_slow_data(),
)
.await;
}
async fn assert_deferred_parity_hedges_slow_data() {
const NUM_SHARDS: usize = 1; const NUM_SHARDS: usize = 1;
const BLOCK_SIZE: usize = 64; const BLOCK_SIZE: usize = 64;
const DATA_SHARDS: usize = 2; const DATA_SHARDS: usize = 2;
@@ -4951,33 +4970,27 @@ mod tests {
]; ];
let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE); let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE);
let (bufs, errs, engaged, readers_remaining) = with_decode_read_policy(DecodeReadPolicy::DemandBound, async { let mut parallel_reader = ParallelReader::new_with_metrics_path_read_costs_timeout_and_reconstruction_verification(
let mut parallel_reader = ParallelReader::new_with_metrics_path_read_costs_timeout_and_reconstruction_verification( readers,
readers, erasure,
erasure, 0,
0, NUM_SHARDS * BLOCK_SIZE,
NUM_SHARDS * BLOCK_SIZE, None,
None, vec![ShardReadCost::Unknown; DATA_SHARDS + PARITY_SHARDS],
vec![ShardReadCost::Unknown; DATA_SHARDS + PARITY_SHARDS], Duration::from_secs(60),
Duration::from_secs(60), true,
true, );
); let (bufs, errs) = tokio::time::timeout(Duration::from_secs(2), parallel_reader.read())
let (bufs, errs) = tokio::time::timeout(Duration::from_secs(2), parallel_reader.read()) .await
.await .expect("deferred parity must cover a hedged data shard without waiting for read_timeout");
.expect("deferred parity must cover a hedged data shard without waiting for read_timeout");
(
bufs,
errs,
parallel_reader.engaged.clone(),
parallel_reader.readers.iter().map(Option::is_some).collect::<Vec<_>>(),
)
})
.await;
assert!(matches!(&errs[0], Some(DiskError::Io(err)) if err.kind() == ErrorKind::TimedOut)); assert!(matches!(&errs[0], Some(DiskError::Io(err)) if err.kind() == ErrorKind::TimedOut));
assert_eq!(bufs.iter().filter(|buf| buf.is_some()).count(), DATA_SHARDS + 1); assert_eq!(bufs.iter().filter(|buf| buf.is_some()).count(), DATA_SHARDS + 1);
assert_eq!(engaged.as_slice(), &[true, true, true, true]); assert_eq!(parallel_reader.engaged.as_slice(), &[true, true, true, true]);
assert_eq!(readers_remaining, vec![false, true, true, true]); assert_eq!(
parallel_reader.readers.iter().map(Option::is_some).collect::<Vec<_>>(),
vec![false, true, true, true]
);
} }
/// A fast data failure must admit deferred parity immediately. There is /// A fast data failure must admit deferred parity immediately. There is
@@ -5046,6 +5059,24 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn test_demand_bound_canceled_hedge_preserves_deferred_parity_for_next_stripe() { async fn test_demand_bound_canceled_hedge_preserves_deferred_parity_for_next_stripe() {
with_decode_read_policy(
DecodeReadPolicy::DemandBound,
assert_canceled_hedge_preserves_deferred_parity_for_next_stripe(),
)
.await;
}
#[tokio::test]
#[serial_test::serial]
async fn test_data_shards_only_gate_canceled_hedge_preserves_deferred_parity_for_next_stripe() {
temp_env::async_with_vars(
[(ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE, Some("true"))],
assert_canceled_hedge_preserves_deferred_parity_for_next_stripe(),
)
.await;
}
async fn assert_canceled_hedge_preserves_deferred_parity_for_next_stripe() {
const BLOCK_SIZE: usize = 64; const BLOCK_SIZE: usize = 64;
const DATA_SHARDS: usize = 2; const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2; const PARITY_SHARDS: usize = 2;
@@ -5094,7 +5125,7 @@ mod tests {
Some(BitrotReader::new(TestShardReader::Pending, SHARD_SIZE, hash_algo, false)), Some(BitrotReader::new(TestShardReader::Pending, SHARD_SIZE, hash_algo, false)),
]; ];
let (first_parity_reserved, second_result) = with_decode_read_policy(DecodeReadPolicy::DemandBound, async { let (first_parity_reserved, second_result) = {
let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE); let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE);
let mut parallel_reader = ParallelReader::new_with_metrics_path_read_timeout_and_reconstruction_verification( let mut parallel_reader = ParallelReader::new_with_metrics_path_read_timeout_and_reconstruction_verification(
readers, readers,
@@ -5155,8 +5186,7 @@ mod tests {
parallel_reader.readers[2].is_some() && parallel_reader.readers[3].is_some(), parallel_reader.readers[2].is_some() && parallel_reader.readers[3].is_some(),
(third_buffers, third_errors), (third_buffers, third_errors),
) )
}) };
.await;
assert!(first_parity_reserved); assert!(first_parity_reserved);
assert_eq!(parity_calls.load(Ordering::SeqCst), PARITY_SHARDS * 2); assert_eq!(parity_calls.load(Ordering::SeqCst), PARITY_SHARDS * 2);
@@ -1974,14 +1974,10 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
return; return;
} }
// Only CopySource uses disposable, stripe-aligned reopeners. Ordinary GET // Every demand-bound lockstep reader needs a disposable, stripe-aligned
// readers use the existing deferred handle and should not retain one // reopener. Otherwise a recovered slow data read can cancel and consume
// heap-allocated closure (plus cloned path/disk state) for every parity // the only parity reserve needed by a later degraded stripe.
// slot. let demand_bound_lockstep = crate::erasure::coding::decode::get_lockstep_data_shards_only_enabled();
let copy_source_demand_bound = matches!(
crate::set_disk::get_object_read_policy(),
crate::set_disk::GetObjectReadPolicy::CopySource
);
for idx in 0..disks.len() { for idx in 0..disks.len() {
if setup.attempted[idx] { if setup.attempted[idx] {
@@ -1996,7 +1992,7 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
let disk = disks[idx].clone(); let disk = disks[idx].clone();
let data_dir = files[idx].data_dir.unwrap_or_default(); let data_dir = files[idx].data_dir.unwrap_or_default();
let path = format!("{object}/{data_dir}/part.{part_number}"); let path = format!("{object}/{data_dir}/part.{part_number}");
let reopener = copy_source_demand_bound.then(|| { let reopener = demand_bound_lockstep.then(|| {
deferred_reader_reopener( deferred_reader_reopener(
inline_data.clone(), inline_data.clone(),
disk.clone(), disk.clone(),
@@ -2037,7 +2033,7 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
// ready/error bookkeeping that quorum decisions rely on is left untouched. // ready/error bookkeeping that quorum decisions rely on is left untouched.
// Gate off (default): keep the eagerly opened parity readers exactly as // Gate off (default): keep the eagerly opened parity readers exactly as
// before — the lockstep path reads them on every stripe. // before — the lockstep path reads them on every stripe.
if !crate::erasure::coding::decode::get_lockstep_data_shards_only_enabled() { if !demand_bound_lockstep {
return; return;
} }
for idx in data_shards..disks.len() { for idx in data_shards..disks.len() {
@@ -2049,7 +2045,7 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
let disk = disks[idx].clone(); let disk = disks[idx].clone();
let data_dir = files[idx].data_dir.unwrap_or_default(); let data_dir = files[idx].data_dir.unwrap_or_default();
let path = format!("{object}/{data_dir}/part.{part_number}"); let path = format!("{object}/{data_dir}/part.{part_number}");
let reopener = copy_source_demand_bound.then(|| { let reopener = demand_bound_lockstep.then(|| {
deferred_reader_reopener( deferred_reader_reopener(
inline_data.clone(), inline_data.clone(),
disk.clone(), disk.clone(),
+21 -9
View File
@@ -5564,9 +5564,10 @@ mod tests {
/// backlog#923: with the data-shards-only lockstep gate on, every retained /// backlog#923: with the data-shards-only lockstep gate on, every retained
/// parity reader must be an unopened deferred reader carrying a stripe /// parity reader must be an unopened deferred reader carrying a stripe
/// handle, so the decode path can realign it to a mid-object stripe. With /// handle and disposable reopener, so the decode path can realign it to a
/// the gate off (default), eagerly opened parity readers are kept exactly /// mid-object stripe without consuming the later-stripe reserve. With the
/// as before and carry no handles. /// gate off (default), eagerly opened parity readers are kept exactly as
/// before and carry neither.
#[tokio::test] #[tokio::test]
#[serial_test::serial] #[serial_test::serial]
async fn bitrot_reader_setup_gates_parity_stripe_handle_conversion() { async fn bitrot_reader_setup_gates_parity_stripe_handle_conversion() {
@@ -5595,6 +5596,11 @@ mod tests {
enabled.is_some(), enabled.is_some(),
"parity slot {idx} stripe handle must match the gate (enabled={enabled:?})" "parity slot {idx} stripe handle must match the gate (enabled={enabled:?})"
); );
assert_eq!(
setup.deferred_reopeners[idx].is_some(),
enabled.is_some(),
"parity slot {idx} reopener must match the gate (enabled={enabled:?})"
);
} }
if enabled.is_some() { if enabled.is_some() {
@@ -5617,13 +5623,17 @@ mod tests {
} }
#[tokio::test] #[tokio::test]
#[serial_test::serial]
async fn bitrot_reader_setup_data_blocks_first_keeps_deferred_fallback_readers() { async fn bitrot_reader_setup_data_blocks_first_keeps_deferred_fallback_readers() {
let mut setup = setup_inline_bitrot_readers_with_env( let mut setup = temp_env::async_with_vars(
vec![Some(b"aaaa"), Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")], [("RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE", Some("true"))],
2, setup_inline_bitrot_readers_with_env(
2, vec![Some(b"aaaa"), Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
BitrotReaderSetupMode::ReadQuorum, 2,
true, 2,
BitrotReaderSetupMode::ReadQuorum,
true,
),
) )
.await; .await;
@@ -5631,6 +5641,8 @@ mod tests {
assert_eq!(setup.available_shards(), 2); assert_eq!(setup.available_shards(), 2);
assert_eq!(setup.scheduled_shards(), 2); assert_eq!(setup.scheduled_shards(), 2);
assert_eq!(setup.readers.iter().filter(|reader| reader.is_some()).count(), 4); assert_eq!(setup.readers.iter().filter(|reader| reader.is_some()).count(), 4);
assert!(setup.deferred_reopeners[2].is_some());
assert!(setup.deferred_reopeners[3].is_some());
let fallback_index = setup let fallback_index = setup
.attempted .attempted