mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-04 04:17:44 +00:00
perf(ecstore): data-shards-only lockstep GET reads with stripe-aligned deferred parity engagement (opt-in) (#4392)
* feat(ecstore): add stripe-advance handles for deferred bitrot readers Give DeferredObjectReader a shared pending state and expose a DeferredReaderStripeHandle that advances the still-unopened source by whole bitrot blocks using the same bitrot_encoded_range geometry the reader was created with (identity mapping when hash_size == 0). This lets the GET decode path open a parity shard aligned to the stripe where a data shard failed instead of reading every parity shard on every stripe (backlog#923). An already-opened (or failed) reader rejects the advance so callers retire it rather than engage it out of alignment; bitrot verification after an advance checks the advanced stripe's block against that stripe's stored hash. Co-Authored-By: heihutu <heihutu@gmail.com> * perf(ecstore): read only data shards on healthy lockstep GET behind opt-in gate PR #4289's lockstep fix made every reconstruction-verifying GET read all data+parity shards per stripe; the parity blocks are read, bitrot-hashed and then discarded, a deterministic 2x read-bytes/IOPS/hash-CPU amplification on healthy 2+2 objects (backlog#923). With the new opt-in gate RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE=true (default: false, behavior identical to main): - read_lockstep keeps only the data slots engaged while the object is healthy; parity slots stay unopened deferred readers. - When a data shard is missing or dies at stripe k, parity readers are engaged mid-object by advancing their deferred stripe handle to stripe k, preserving the lockstep alignment invariant from backlog#832. - Degraded stripes engage one parity beyond the decode quorum so reconstruction verification keeps an extra source to check against (erasure.rs only verifies when available > data shards); an engaged parity reader that errors is retired for the rest of the object like any other, and a parity reader that cannot be realigned is retired instead of being read out of position. - fill_deferred_bitrot_readers records stripe handles for deferred slots and, gate-on only, swaps eagerly opened parity readers for unopened deferred ones so they remain engageable mid-object; ready/error bookkeeping used by quorum decisions is untouched. - Both GET paths (legacy duplex via Erasure::decode_with_stripe_handles, codec streaming via ParallelReader::with_deferred_parity_handles) carry the handles from reader setup. Short-read -> UnexpectedEof -> whole-object retirement and the inconsistent-source rejection are unchanged in both gate modes; tests lock the healthy-path data-shards-only call counts, the default read-all-shards behavior, mid-object parity engagement for streaming and hash_size==0 formats, and mid-stream inconsistent-parity rejection. Co-Authored-By: heihutu <heihutu@gmail.com> --------- Co-authored-by: heihutu <heihutu@gmail.com>
This commit is contained in:
@@ -24,6 +24,7 @@ use crate::disk::error::Error;
|
||||
use crate::disk::error_reduce::reduce_errs;
|
||||
use crate::erasure::codec::workspace::ShardBufferPool;
|
||||
use crate::erasure::coding::{BitrotReader, Erasure};
|
||||
use crate::io_support::bitrot::DeferredReaderStripeHandle;
|
||||
use crate::set_disk::shard_source::{ShardReadCost, ShardStripeSource, StripeReadState};
|
||||
use futures::FutureExt;
|
||||
use futures::stream::{FuturesUnordered, StreamExt};
|
||||
@@ -112,6 +113,26 @@ fn get_decode_stripe_prefetch_count() -> usize {
|
||||
const ENV_RUSTFS_GET_BITROT_DECODE_OVERLAP_ENABLE: &str = "RUSTFS_GET_BITROT_DECODE_OVERLAP_ENABLE";
|
||||
const DEFAULT_RUSTFS_GET_BITROT_DECODE_OVERLAP_ENABLE: bool = false;
|
||||
|
||||
/// Enable the data-shards-only lockstep GET read (backlog#923).
|
||||
/// When enabled, the reconstruction-verifying (lockstep) GET path reads only
|
||||
/// the data shards per stripe while the object is healthy; parity slots stay
|
||||
/// unopened deferred readers that are engaged — realigned to the failing
|
||||
/// stripe via their stripe handle — only when a data shard is missing or dies.
|
||||
/// This halves per-GET read bytes, IOPS, and bitrot hashing on healthy
|
||||
/// objects with 2+2 layouts.
|
||||
/// Default: false (current behavior: every live shard reader is read on every
|
||||
/// stripe).
|
||||
const ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE: &str = "RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE";
|
||||
const DEFAULT_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE: bool = false;
|
||||
|
||||
/// Whether the data-shards-only lockstep GET read is enabled (backlog#923).
|
||||
pub(crate) fn get_lockstep_data_shards_only_enabled() -> bool {
|
||||
rustfs_utils::get_env_bool(
|
||||
ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE,
|
||||
DEFAULT_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE,
|
||||
)
|
||||
}
|
||||
|
||||
/// Get whether bitrot-decode overlap is enabled.
|
||||
fn is_bitrot_decode_overlap_enabled() -> bool {
|
||||
rustfs_utils::get_env_bool(
|
||||
@@ -327,6 +348,14 @@ pub(crate) struct ParallelReader<R> {
|
||||
// Request-scoped shard buffers keyed by shard index. Keeping ownership in
|
||||
// `ParallelReader` avoids dropping unused parity/backup slot buffers between stripes.
|
||||
buffers: ShardBufferPool,
|
||||
// Lockstep-path state (verify_reconstruction == true). `engaged[i]` marks
|
||||
// readers that participate in each stripe read: all data slots from the
|
||||
// start, parity slots only once a data shard is missing/dead. Unengaged
|
||||
// parity stays an unopened deferred reader; `deferred_handles[i]` realigns
|
||||
// it to the current stripe when it is engaged mid-object (backlog#923).
|
||||
engaged: Vec<bool>,
|
||||
deferred_handles: Vec<Option<DeferredReaderStripeHandle>>,
|
||||
stripe_index: usize,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -502,6 +531,14 @@ where
|
||||
|
||||
// Ensure offset does not exceed shard_file_size
|
||||
|
||||
// Default (gate off): every slot is engaged, i.e. the lockstep path
|
||||
// reads all live readers on every stripe — the pre-backlog#923
|
||||
// behavior. With the gate on, only data slots start engaged; parity is
|
||||
// engaged on demand, stripe-aligned through its deferred handle.
|
||||
let data_shards_only = get_lockstep_data_shards_only_enabled();
|
||||
let engaged = (0..readers.len())
|
||||
.map(|index| !data_shards_only || index < e.data_shards)
|
||||
.collect();
|
||||
ParallelReader {
|
||||
readers,
|
||||
offset,
|
||||
@@ -515,8 +552,23 @@ where
|
||||
verify_reconstruction,
|
||||
locality_preference_enabled: get_shard_locality_preference_enabled(),
|
||||
buffers: ShardBufferPool::new(e.data_shards + e.parity_shards),
|
||||
engaged,
|
||||
deferred_handles: Vec::new(),
|
||||
stripe_index: 0,
|
||||
}
|
||||
}
|
||||
|
||||
/// Attach the per-slot deferred stripe handles produced during bitrot
|
||||
/// reader setup. Only parity slots are consulted: a handle lets the
|
||||
/// lockstep path open a parity shard aligned to the stripe where a data
|
||||
/// shard failed, instead of reading every parity shard on every stripe
|
||||
/// (backlog#923).
|
||||
pub(crate) fn with_deferred_parity_handles(mut self, mut handles: Vec<Option<DeferredReaderStripeHandle>>) -> Self {
|
||||
handles.resize_with(self.readers.len(), || None);
|
||||
handles.truncate(self.readers.len());
|
||||
self.deferred_handles = handles;
|
||||
self
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
@@ -983,18 +1035,29 @@ where
|
||||
|
||||
/// Lockstep stripe read for the reconstruction-verifying GET path.
|
||||
///
|
||||
/// Reads every still-live shard reader exactly once per stripe and waits for
|
||||
/// all of them, so all readers advance by one block per stripe and remain
|
||||
/// mutually block-aligned. This removes the desync that produced
|
||||
/// "inconsistent read source shards" (backlog#832): with the adaptive
|
||||
/// data-first path a parity reader pulled in as a substitute mid-object was
|
||||
/// still at its stream start and returned an earlier stripe than the data
|
||||
/// shards, so reconstruction verification correctly rejected the mismatch
|
||||
/// and the large-object GET truncated mid-stream.
|
||||
/// Reads every *engaged* shard reader exactly once per stripe and waits for
|
||||
/// all of them, so all engaged readers advance by one block per stripe and
|
||||
/// remain mutually block-aligned. This preserves the alignment invariant
|
||||
/// that removed the "inconsistent read source shards" desync (backlog#832):
|
||||
/// a reader must never contribute a block from a different stripe than the
|
||||
/// rest of the set.
|
||||
///
|
||||
/// Any reader that errors is retired for the rest of the object: a streaming
|
||||
/// shard read that failed mid-block can no longer be trusted to be aligned,
|
||||
/// so re-reading it on a later stripe would reintroduce the desync.
|
||||
/// While the object is healthy only the data shards are engaged, so a
|
||||
/// healthy GET reads exactly `data_shards` shards per stripe instead of all
|
||||
/// `data + parity` shards (backlog#923). Parity slots stay unopened
|
||||
/// deferred readers; when a data shard is missing or dies at stripe `k`,
|
||||
/// enough parity readers are engaged — realigned to stripe `k` through
|
||||
/// their [`DeferredReaderStripeHandle`] — to restore the decode quorum
|
||||
/// *plus one extra shard* so reconstruction verification still has a
|
||||
/// source to check against (`erasure.rs` only verifies when
|
||||
/// `available_shards > data_shards`).
|
||||
///
|
||||
/// Any reader that errors is retired for the rest of the object — data and
|
||||
/// newly engaged parity alike: a streaming shard read that failed mid-block
|
||||
/// can no longer be trusted to be aligned, so re-reading it on a later
|
||||
/// stripe would reintroduce the desync. A parity reader that cannot be
|
||||
/// realigned (no pending deferred handle) is likewise retired instead of
|
||||
/// being read out of position.
|
||||
async fn read_lockstep(&mut self) -> (Vec<Option<Vec<u8>>>, Vec<Option<Error>>) {
|
||||
let num_readers = self.readers.len();
|
||||
let shard_size = if self.offset + self.shard_size > self.shard_file_size {
|
||||
@@ -1012,14 +1075,42 @@ where
|
||||
// Advance to the next stripe (see the matching note in `read`); the
|
||||
// lockstep path must track stripe geometry identically (backlog#799 B2).
|
||||
self.offset += shard_size;
|
||||
let stripe_index = self.stripe_index;
|
||||
self.stripe_index += 1;
|
||||
|
||||
self.buffers.ensure_slots(num_readers);
|
||||
|
||||
// Engage parity up front when data slots are already known to be
|
||||
// missing (offline at setup or retired on an earlier stripe), so the
|
||||
// substitute reads run in parallel with the surviving data reads.
|
||||
let missing_data_readers = self.readers.iter().take(self.data_shards).filter(|r| r.is_none()).count();
|
||||
if missing_data_readers > 0 {
|
||||
// One extra engaged parity beyond the reconstruction quorum keeps
|
||||
// reconstruction verification effective: with exactly
|
||||
// `data_shards` available shards the verification would silently
|
||||
// turn itself off.
|
||||
let want = missing_data_readers + 1;
|
||||
let mut have = (self.data_shards..num_readers)
|
||||
.filter(|&i| self.engaged[i] && self.readers[i].is_some())
|
||||
.count();
|
||||
for idx in self.data_shards..num_readers {
|
||||
if have >= want {
|
||||
break;
|
||||
}
|
||||
if self.readers[idx].is_some() && !self.engaged[idx] && self.try_engage_parity(idx, stripe_index) {
|
||||
have += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pre-claim per-slot buffers so the `self.readers` borrow below stays
|
||||
// disjoint from `self.buffers`.
|
||||
let has_reader: Vec<bool> = self.readers.iter().map(Option::is_some).collect();
|
||||
let participating: Vec<bool> = (0..num_readers)
|
||||
.map(|i| self.engaged[i] && self.readers[i].is_some())
|
||||
.collect();
|
||||
let mut bufs: Vec<Option<Vec<u8>>> = Vec::with_capacity(num_readers);
|
||||
for (i, present) in has_reader.iter().enumerate() {
|
||||
bufs.push(if *present {
|
||||
for (i, participates) in participating.iter().enumerate() {
|
||||
bufs.push(if *participates {
|
||||
Some(self.buffers.take(i, shard_size))
|
||||
} else {
|
||||
None
|
||||
@@ -1045,7 +1136,7 @@ where
|
||||
let mut sets = FuturesUnordered::new();
|
||||
let reader_iter = ReaderLaunchIter::new(&mut self.readers, &read_costs, locality_preference_enabled);
|
||||
for (i, reader) in reader_iter {
|
||||
if reader.is_none() {
|
||||
if reader.is_none() || !participating[i] {
|
||||
continue;
|
||||
}
|
||||
let read_cost = read_costs.get(i).copied().unwrap_or(ShardReadCost::Unknown);
|
||||
@@ -1085,6 +1176,51 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
// A data shard may have died during this stripe's reads. The unengaged
|
||||
// parity readers are still unopened, so they can be aligned to *this*
|
||||
// stripe and read now: engage them one at a time until the stripe has
|
||||
// `data_shards + 1` successful shards (decode quorum plus one
|
||||
// reconstruction-verification source) or parity runs out.
|
||||
loop {
|
||||
let success_now = shards.iter().filter(|shard| shard.is_some()).count();
|
||||
let data_shard_missing = shards.iter().take(data_shards).any(|shard| shard.is_none());
|
||||
if !data_shard_missing || success_now > data_shards {
|
||||
break;
|
||||
}
|
||||
let Some(idx) = (data_shards..num_readers).find(|&i| self.readers[i].is_some() && !self.engaged[i]) else {
|
||||
break;
|
||||
};
|
||||
if !self.try_engage_parity(idx, stripe_index) {
|
||||
continue;
|
||||
}
|
||||
let read_cost = read_costs.get(idx).copied().unwrap_or(ShardReadCost::Unknown);
|
||||
let recycled_buf = Some(self.buffers.take(idx, shard_size));
|
||||
scheduled += 1;
|
||||
let (i, _read_cost, result, _should_retire) = read_shard(
|
||||
idx,
|
||||
read_cost,
|
||||
&mut self.readers[idx],
|
||||
recycled_buf,
|
||||
shard_size,
|
||||
data_shards,
|
||||
read_timeout,
|
||||
metrics_path,
|
||||
)
|
||||
.await;
|
||||
completed += 1;
|
||||
match result {
|
||||
Ok(v) => {
|
||||
shards[i] = Some(v);
|
||||
success += 1;
|
||||
}
|
||||
Err(e) => {
|
||||
errs[i] = Some(e);
|
||||
retire_readers.push(i);
|
||||
failed += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if let Some(path) = metrics_path {
|
||||
record_get_stage_duration_if_enabled(path, GET_STAGE_STRIPE_READ_QUORUM, stripe_read_start);
|
||||
rustfs_io_metrics::record_get_object_shard_read_fanout(path, scheduled, completed, success, failed);
|
||||
@@ -1097,6 +1233,40 @@ where
|
||||
(shards, errs)
|
||||
}
|
||||
|
||||
/// Attempt to bring an as-yet-unread parity reader into the lockstep read
|
||||
/// set at `stripe_index`.
|
||||
///
|
||||
/// At stripe 0 every reader is still positioned at the stream start, so
|
||||
/// engagement is trivially aligned. Past stripe 0 the parity reader must
|
||||
/// still be an unopened deferred reader: its pending open offset is
|
||||
/// advanced by `stripe_index` bitrot blocks (the `bitrot_encoded_range`
|
||||
/// geometry) so its first read returns the current stripe. A parity reader
|
||||
/// that cannot be realigned is retired for the rest of the object,
|
||||
/// mirroring the retire-on-error rule: reading it would return an earlier
|
||||
/// stripe and reintroduce the backlog#832 desync.
|
||||
fn try_engage_parity(&mut self, idx: usize, stripe_index: usize) -> bool {
|
||||
if stripe_index == 0 {
|
||||
self.engaged[idx] = true;
|
||||
return true;
|
||||
}
|
||||
let advanced = self
|
||||
.deferred_handles
|
||||
.get(idx)
|
||||
.and_then(|handle| handle.as_ref())
|
||||
.is_some_and(|handle| handle.advance_stripes(stripe_index));
|
||||
if advanced {
|
||||
self.engaged[idx] = true;
|
||||
true
|
||||
} else {
|
||||
warn!(
|
||||
shard_index = idx,
|
||||
stripe_index, "retiring parity reader that cannot be aligned to the current stripe"
|
||||
);
|
||||
self.readers[idx] = None;
|
||||
false
|
||||
}
|
||||
}
|
||||
|
||||
pub fn recycle_shards(&mut self, shards: &mut [Option<Vec<u8>>]) {
|
||||
for (i, reader) in self.readers.iter().enumerate() {
|
||||
if reader.is_some()
|
||||
@@ -1274,7 +1444,8 @@ impl Erasure {
|
||||
W: AsyncWrite + Send + Sync + Unpin,
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
self.decode_inner(writer, readers, offset, length, total_length, None).await
|
||||
self.decode_inner(writer, readers, offset, length, total_length, None, Vec::new())
|
||||
.await
|
||||
}
|
||||
|
||||
pub(crate) async fn decode_with_read_costs<W, R>(
|
||||
@@ -1290,10 +1461,33 @@ impl Erasure {
|
||||
W: AsyncWrite + Send + Sync + Unpin,
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
self.decode_inner(writer, readers, offset, length, total_length, Some(read_costs))
|
||||
self.decode_inner(writer, readers, offset, length, total_length, Some(read_costs), Vec::new())
|
||||
.await
|
||||
}
|
||||
|
||||
/// GET decode entry point that also carries the deferred-parity stripe
|
||||
/// handles from bitrot reader setup, so unengaged parity readers can be
|
||||
/// opened aligned to the stripe where a data shard fails (backlog#923).
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(crate) async fn decode_with_stripe_handles<W, R>(
|
||||
&self,
|
||||
writer: &mut W,
|
||||
readers: Vec<Option<BitrotReader<R>>>,
|
||||
offset: usize,
|
||||
length: usize,
|
||||
total_length: usize,
|
||||
read_costs: Option<Vec<ShardReadCost>>,
|
||||
deferred_handles: Vec<Option<DeferredReaderStripeHandle>>,
|
||||
) -> (usize, Option<std::io::Error>)
|
||||
where
|
||||
W: AsyncWrite + Send + Sync + Unpin,
|
||||
R: AsyncRead + Unpin + Send + Sync,
|
||||
{
|
||||
self.decode_inner(writer, readers, offset, length, total_length, read_costs, deferred_handles)
|
||||
.await
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
async fn decode_inner<W, R>(
|
||||
&self,
|
||||
writer: &mut W,
|
||||
@@ -1302,6 +1496,7 @@ impl Erasure {
|
||||
length: usize,
|
||||
total_length: usize,
|
||||
read_costs: Option<Vec<ShardReadCost>>,
|
||||
deferred_handles: Vec<Option<DeferredReaderStripeHandle>>,
|
||||
) -> (usize, Option<std::io::Error>)
|
||||
where
|
||||
W: AsyncWrite + Send + Sync + Unpin,
|
||||
@@ -1349,7 +1544,8 @@ impl Erasure {
|
||||
)
|
||||
} else {
|
||||
ParallelReader::new_for_decode(readers, self.clone(), offset, total_length, Some(GET_OBJECT_PATH_LEGACY_DUPLEX))
|
||||
};
|
||||
}
|
||||
.with_deferred_parity_handles(deferred_handles);
|
||||
|
||||
let start = offset / self.block_size;
|
||||
let end = end_offset.saturating_sub(1) / self.block_size;
|
||||
@@ -1479,9 +1675,82 @@ mod tests {
|
||||
use rustfs_utils::HashAlgorithm;
|
||||
use std::io::Cursor;
|
||||
use std::pin::Pin;
|
||||
use std::sync::{
|
||||
Arc,
|
||||
atomic::{AtomicUsize, Ordering},
|
||||
};
|
||||
use std::task::{Context, Poll};
|
||||
use tokio::io::ReadBuf;
|
||||
|
||||
type BoxedShardReader = Box<dyn AsyncRead + Send + Sync + Unpin>;
|
||||
|
||||
/// Counts the raw bytes pulled from a shard stream, to prove which shards
|
||||
/// a decode path actually touches (backlog#923 call-count evidence).
|
||||
struct CountingShardReader {
|
||||
inner: Cursor<Vec<u8>>,
|
||||
bytes_read: Arc<AtomicUsize>,
|
||||
}
|
||||
|
||||
impl AsyncRead for CountingShardReader {
|
||||
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<io::Result<()>> {
|
||||
let before = buf.filled().len();
|
||||
let result = Pin::new(&mut self.inner).poll_read(cx, buf);
|
||||
if let Poll::Ready(Ok(())) = result {
|
||||
let delta = buf.filled().len() - before;
|
||||
self.bytes_read.fetch_add(delta, Ordering::SeqCst);
|
||||
}
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
/// Build a production-shaped reader set: data shards as opened stream
|
||||
/// readers, parity shards as *unopened* deferred readers carrying stripe
|
||||
/// handles (mirroring `fill_deferred_bitrot_readers`). `truncate` shortens
|
||||
/// the given data shard buffers to simulate a reader dying mid-object.
|
||||
fn readers_with_deferred_parity(
|
||||
shard_bufs: &[Vec<u8>],
|
||||
data_shards: usize,
|
||||
shard_size: usize,
|
||||
hash_algo: &HashAlgorithm,
|
||||
truncate: &[(usize, usize)],
|
||||
) -> (Vec<Option<BitrotReader<BoxedShardReader>>>, Vec<Option<DeferredReaderStripeHandle>>) {
|
||||
use crate::io_support::bitrot::create_deferred_bitrot_reader_with_stripe_handle;
|
||||
|
||||
let mut readers = Vec::with_capacity(shard_bufs.len());
|
||||
let mut handles: Vec<Option<DeferredReaderStripeHandle>> = vec![None; shard_bufs.len()];
|
||||
for (i, buf) in shard_bufs.iter().enumerate() {
|
||||
if i < data_shards {
|
||||
let bytes = truncate
|
||||
.iter()
|
||||
.find(|(index, _)| *index == i)
|
||||
.map(|(_, len)| buf[..*len].to_vec())
|
||||
.unwrap_or_else(|| buf.clone());
|
||||
readers.push(Some(BitrotReader::new(
|
||||
Box::new(Cursor::new(bytes)) as BoxedShardReader,
|
||||
shard_size,
|
||||
hash_algo.clone(),
|
||||
false,
|
||||
)));
|
||||
} else {
|
||||
let (reader, handle) = create_deferred_bitrot_reader_with_stripe_handle(
|
||||
Some(bytes::Bytes::from(buf.clone())),
|
||||
None,
|
||||
"test-bucket",
|
||||
"test-object",
|
||||
0,
|
||||
buf.len(),
|
||||
shard_size,
|
||||
hash_algo.clone(),
|
||||
false,
|
||||
false,
|
||||
);
|
||||
readers.push(Some(reader));
|
||||
handles[i] = Some(handle);
|
||||
}
|
||||
}
|
||||
(readers, handles)
|
||||
}
|
||||
|
||||
enum TestShardReader {
|
||||
Ready(Cursor<Vec<u8>>),
|
||||
Pending,
|
||||
@@ -1722,17 +1991,9 @@ mod tests {
|
||||
assert!(output.is_empty());
|
||||
}
|
||||
|
||||
/// Regression for backlog#832: a data shard reader that dies partway through
|
||||
/// a multi-stripe object must still reconstruct byte-exact output. The parity
|
||||
/// readers that fill in for the failed data shard must be aligned to the
|
||||
/// current stripe. Before the lockstep read path, parity readers were only
|
||||
/// pulled in on demand and were still positioned at their stream start
|
||||
/// (block 0), so they returned an earlier stripe than the surviving data
|
||||
/// shards — every shard passed its own bitrot hash yet the set was mutually
|
||||
/// inconsistent, so reconstruction verification rejected it ("inconsistent
|
||||
/// read source shards") and the large-object GET truncated mid-stream.
|
||||
#[tokio::test]
|
||||
async fn test_erasure_decode_recovers_when_data_shard_dies_midway() {
|
||||
/// One mid-object data-shard death scenario, shared by the gate-on and
|
||||
/// gate-off variants of `test_erasure_decode_recovers_when_data_shard_dies_midway`.
|
||||
async fn run_decode_midway_death_case(hash_algo: HashAlgorithm, context: &str) {
|
||||
const DATA_SHARDS: usize = 4;
|
||||
const PARITY_SHARDS: usize = 2;
|
||||
const BLOCK_SIZE: usize = 64;
|
||||
@@ -1745,7 +2006,6 @@ mod tests {
|
||||
let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE);
|
||||
let total_shards = DATA_SHARDS + PARITY_SHARDS;
|
||||
let shard_size = erasure.shard_size();
|
||||
let hash_algo = HashAlgorithm::HighwayHash256;
|
||||
let hash_size = hash_algo.size();
|
||||
|
||||
let mut shard_writers: Vec<BitrotWriter<Cursor<Vec<u8>>>> = (0..total_shards)
|
||||
@@ -1764,29 +2024,228 @@ mod tests {
|
||||
|
||||
let shard_bufs: Vec<Vec<u8>> = shard_writers.into_iter().map(|w| w.into_inner().into_inner()).collect();
|
||||
|
||||
// Truncate data shard 0 to only its first (hash+data) block: it reads the
|
||||
// first stripe successfully, then errors (UnexpectedEof reading the next
|
||||
// block's hash), forcing every later stripe to be reconstructed from the
|
||||
// parity shards. With all readers present the reconstruction must succeed.
|
||||
// Truncate data shard 0 to only its first (hash+data) block: it reads
|
||||
// the first stripe successfully, then errors (UnexpectedEof), forcing
|
||||
// every later stripe to engage the deferred parity readers aligned to
|
||||
// the failing stripe.
|
||||
let first_block_len = (hash_size + shard_size).min(shard_bufs[0].len());
|
||||
let readers: Vec<Option<BitrotReader<Cursor<Vec<u8>>>>> = shard_bufs
|
||||
let (readers, handles) =
|
||||
readers_with_deferred_parity(&shard_bufs, DATA_SHARDS, shard_size, &hash_algo, &[(0, first_block_len)]);
|
||||
|
||||
let mut output = Vec::new();
|
||||
let (written, err) = erasure
|
||||
.decode_with_stripe_handles(&mut output, readers, 0, total_len, total_len, None, handles)
|
||||
.await;
|
||||
assert!(
|
||||
err.is_none(),
|
||||
"{context}, algo={hash_algo:?}: mid-object data-shard failure must still reconstruct: {err:?}"
|
||||
);
|
||||
assert_eq!(
|
||||
written, total_len,
|
||||
"{context}, algo={hash_algo:?}: short write after mid-object shard failure"
|
||||
);
|
||||
assert_eq!(
|
||||
output, total_data,
|
||||
"{context}, algo={hash_algo:?}: reconstructed bytes mismatch (stripe desync?)"
|
||||
);
|
||||
}
|
||||
|
||||
/// Regression for backlog#832 (extended for backlog#923): a data shard
|
||||
/// reader that dies partway through a multi-stripe object must still
|
||||
/// reconstruct byte-exact output, and the parity readers that fill in for
|
||||
/// it must be aligned to the current stripe. With the data-shards-only
|
||||
/// gate on, parity slots are unopened deferred readers whose stripe handle
|
||||
/// advances the pending open offset to the failing stripe — mirroring the
|
||||
/// production GET reader setup. Covers both the streaming per-block
|
||||
/// checksum format and the no-per-block-hash (`hash_size == 0`) format,
|
||||
/// where the offset mapping degrades to the identity function (merge gate
|
||||
/// for backlog#923). Also run with the gate off to lock the default
|
||||
/// (read-all-shards) behavior.
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn test_erasure_decode_recovers_when_data_shard_dies_midway() {
|
||||
for enabled in [None, Some("true")] {
|
||||
temp_env::async_with_vars([(ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE, enabled)], async {
|
||||
let context = if enabled.is_some() { "data-shards-only" } else { "default" };
|
||||
for hash_algo in [HashAlgorithm::HighwayHash256, HashAlgorithm::None] {
|
||||
run_decode_midway_death_case(hash_algo, context).await;
|
||||
}
|
||||
})
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
/// Decode a healthy 2+2 multi-stripe object through the lockstep GET path
|
||||
/// and return the raw bytes pulled from each shard stream.
|
||||
async fn healthy_lockstep_shard_bytes() -> (usize, Vec<usize>) {
|
||||
const DATA_SHARDS: usize = 2;
|
||||
const PARITY_SHARDS: usize = 2;
|
||||
const BLOCK_SIZE: usize = 64;
|
||||
|
||||
// 200 bytes => 3 full stripes + 1 partial.
|
||||
let total_data: Vec<u8> = (0..200u32).map(|i| i as u8).collect();
|
||||
let total_len = total_data.len();
|
||||
|
||||
let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE);
|
||||
let total_shards = DATA_SHARDS + PARITY_SHARDS;
|
||||
let shard_size = erasure.shard_size();
|
||||
let hash_algo = HashAlgorithm::HighwayHash256;
|
||||
|
||||
let mut shard_writers: Vec<BitrotWriter<Cursor<Vec<u8>>>> = (0..total_shards)
|
||||
.map(|_| BitrotWriter::new(Cursor::new(Vec::new()), shard_size, hash_algo.clone()))
|
||||
.collect();
|
||||
|
||||
let mut offset = 0;
|
||||
while offset < total_len {
|
||||
let end = (offset + BLOCK_SIZE).min(total_len);
|
||||
let shards = erasure.encode_data(&total_data[offset..end]).unwrap();
|
||||
for (i, shard) in shards.iter().enumerate() {
|
||||
shard_writers[i].write(shard).await.unwrap();
|
||||
}
|
||||
offset = end;
|
||||
}
|
||||
|
||||
let shard_bufs: Vec<Vec<u8>> = shard_writers.into_iter().map(|w| w.into_inner().into_inner()).collect();
|
||||
|
||||
let counters: Vec<Arc<AtomicUsize>> = (0..total_shards).map(|_| Arc::new(AtomicUsize::new(0))).collect();
|
||||
let readers: Vec<Option<BitrotReader<CountingShardReader>>> = shard_bufs
|
||||
.iter()
|
||||
.enumerate()
|
||||
.map(|(i, buf)| {
|
||||
let bytes = if i == 0 {
|
||||
buf[..first_block_len].to_vec()
|
||||
} else {
|
||||
buf.clone()
|
||||
};
|
||||
Some(BitrotReader::new(Cursor::new(bytes), shard_size, hash_algo.clone(), false))
|
||||
Some(BitrotReader::new(
|
||||
CountingShardReader {
|
||||
inner: Cursor::new(buf.clone()),
|
||||
bytes_read: Arc::clone(&counters[i]),
|
||||
},
|
||||
shard_size,
|
||||
hash_algo.clone(),
|
||||
false,
|
||||
))
|
||||
})
|
||||
.collect();
|
||||
|
||||
// `Erasure::decode` is the reconstruction-verifying (lockstep) GET path.
|
||||
let mut output = Vec::new();
|
||||
let (written, err) = erasure.decode(&mut output, readers, 0, total_len, total_len).await;
|
||||
assert!(err.is_none(), "mid-object data-shard failure must still reconstruct: {err:?}");
|
||||
assert_eq!(written, total_len, "short write after mid-object shard failure");
|
||||
assert_eq!(output, total_data, "reconstructed bytes mismatch (stripe desync?)");
|
||||
assert!(err.is_none(), "healthy decode must succeed: {err:?}");
|
||||
assert_eq!(written, total_len);
|
||||
assert_eq!(output, total_data);
|
||||
|
||||
(DATA_SHARDS, counters.iter().map(|counter| counter.load(Ordering::SeqCst)).collect())
|
||||
}
|
||||
|
||||
/// Merge gate for backlog#923: with the gate on, the healthy lockstep GET
|
||||
/// path must read only the data shards — parity readers stay unopened and
|
||||
/// contribute zero read bytes. This is the call-count proof of the 2x read
|
||||
/// amplification fix (2+2 layout reads 2 shards per stripe, not 4).
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn test_lockstep_healthy_get_reads_only_data_shards() {
|
||||
temp_env::async_with_vars([(ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE, Some("true"))], async {
|
||||
let (data_shards, shard_bytes) = healthy_lockstep_shard_bytes().await;
|
||||
for (i, bytes) in shard_bytes.iter().enumerate() {
|
||||
if i < data_shards {
|
||||
assert!(*bytes > 0, "data shard {i} must be read on the healthy path");
|
||||
} else {
|
||||
assert_eq!(
|
||||
*bytes, 0,
|
||||
"healthy GET must not read parity shard {i} (2x read amplification, backlog#923)"
|
||||
);
|
||||
}
|
||||
}
|
||||
})
|
||||
.await;
|
||||
}
|
||||
|
||||
/// Compatibility lock: with the gate off (default), the lockstep path
|
||||
/// keeps the pre-backlog#923 behavior and reads every live shard —
|
||||
/// including parity — on every stripe.
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn test_lockstep_default_reads_all_shards_per_stripe() {
|
||||
temp_env::async_with_vars([(ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE, None::<&str>)], async {
|
||||
let (_data_shards, shard_bytes) = healthy_lockstep_shard_bytes().await;
|
||||
for (i, bytes) in shard_bytes.iter().enumerate() {
|
||||
assert!(*bytes > 0, "default lockstep behavior must read shard {i} on every stripe");
|
||||
}
|
||||
})
|
||||
.await;
|
||||
}
|
||||
|
||||
/// Merge gate for backlog#923: reconstruction verification must stay
|
||||
/// active when parity is engaged mid-object. When a data shard dies at
|
||||
/// stripe k, the lockstep path engages one parity reader beyond the decode
|
||||
/// quorum (`erasure.rs` only verifies when `available > data_shards`), so
|
||||
/// a parity shard whose content is inconsistent with the surviving data —
|
||||
/// while still passing its own bitrot hash — must be detected and fail the
|
||||
/// read instead of silently corrupting the reconstructed output. Run with
|
||||
/// the gate on and off: the rejection must hold in both modes.
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn test_erasure_decode_rejects_inconsistent_parity_engaged_midstream() {
|
||||
for enabled in [None, Some("true")] {
|
||||
temp_env::async_with_vars([(ENV_RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE, enabled)], async {
|
||||
run_inconsistent_parity_midstream_case(if enabled.is_some() { "data-shards-only" } else { "default" }).await;
|
||||
})
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
async fn run_inconsistent_parity_midstream_case(context: &str) {
|
||||
const DATA_SHARDS: usize = 2;
|
||||
const PARITY_SHARDS: usize = 2;
|
||||
const BLOCK_SIZE: usize = 64;
|
||||
|
||||
let total_data: Vec<u8> = (0..200u32).map(|i| i as u8).collect();
|
||||
let total_len = total_data.len();
|
||||
|
||||
let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE);
|
||||
let total_shards = DATA_SHARDS + PARITY_SHARDS;
|
||||
let shard_size = erasure.shard_size();
|
||||
let hash_algo = HashAlgorithm::HighwayHash256;
|
||||
let hash_size = hash_algo.size();
|
||||
|
||||
let mut shard_writers: Vec<BitrotWriter<Cursor<Vec<u8>>>> = (0..total_shards)
|
||||
.map(|_| BitrotWriter::new(Cursor::new(Vec::new()), shard_size, hash_algo.clone()))
|
||||
.collect();
|
||||
|
||||
let mut offset = 0;
|
||||
let mut stripe = 0usize;
|
||||
while offset < total_len {
|
||||
let end = (offset + BLOCK_SIZE).min(total_len);
|
||||
let shards = erasure.encode_data(&total_data[offset..end]).unwrap();
|
||||
for (i, shard) in shards.iter().enumerate() {
|
||||
// Corrupt the first parity shard's payload for every stripe
|
||||
// after the first, *before* it is bitrot-hashed: the shard
|
||||
// passes its own hash check but is inconsistent with the
|
||||
// erasure-coded set.
|
||||
if i == DATA_SHARDS && stripe >= 1 {
|
||||
let mut corrupted = shard.to_vec();
|
||||
corrupted[0] ^= 0x80;
|
||||
shard_writers[i].write(&corrupted).await.unwrap();
|
||||
} else {
|
||||
shard_writers[i].write(shard).await.unwrap();
|
||||
}
|
||||
}
|
||||
offset = end;
|
||||
stripe += 1;
|
||||
}
|
||||
|
||||
let shard_bufs: Vec<Vec<u8>> = shard_writers.into_iter().map(|w| w.into_inner().into_inner()).collect();
|
||||
|
||||
// Data shard 0 dies after stripe 0, forcing parity engagement at stripe 1.
|
||||
let first_block_len = (hash_size + shard_size).min(shard_bufs[0].len());
|
||||
let (readers, handles) =
|
||||
readers_with_deferred_parity(&shard_bufs, DATA_SHARDS, shard_size, &hash_algo, &[(0, first_block_len)]);
|
||||
|
||||
let mut output = Vec::new();
|
||||
let (_written, err) = erasure
|
||||
.decode_with_stripe_handles(&mut output, readers, 0, total_len, total_len, None, handles)
|
||||
.await;
|
||||
|
||||
let err = err.unwrap_or_else(|| panic!("{context}: inconsistent parity engaged mid-stream must fail the read"));
|
||||
assert_eq!(err.kind(), ErrorKind::InvalidData, "{context}");
|
||||
assert!(err.to_string().contains("inconsistent read source shards"), "{context}: {err}");
|
||||
}
|
||||
|
||||
#[cfg(feature = "rio-v2")]
|
||||
|
||||
@@ -27,7 +27,7 @@ use rustfs_utils::HashAlgorithm;
|
||||
use std::future::Future;
|
||||
use std::io::{self, Cursor};
|
||||
use std::pin::Pin;
|
||||
use std::sync::Mutex;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::task::{Context, Poll};
|
||||
use std::time::Instant;
|
||||
use tokio::io::{AsyncRead, ReadBuf};
|
||||
@@ -134,7 +134,7 @@ impl AsyncRead for FirstReadMetricsReader {
|
||||
}
|
||||
|
||||
struct DeferredObjectReader {
|
||||
state: Mutex<DeferredObjectReaderState>,
|
||||
state: Arc<Mutex<DeferredObjectReaderState>>,
|
||||
}
|
||||
|
||||
enum DeferredObjectReaderState {
|
||||
@@ -147,7 +147,64 @@ enum DeferredObjectReaderState {
|
||||
impl DeferredObjectReader {
|
||||
fn new(source: BitrotReaderSource) -> Self {
|
||||
Self {
|
||||
state: Mutex::new(DeferredObjectReaderState::Pending(Some(source))),
|
||||
state: Arc::new(Mutex::new(DeferredObjectReaderState::Pending(Some(source)))),
|
||||
}
|
||||
}
|
||||
|
||||
fn stripe_handle(&self, stripe_stride: usize) -> DeferredReaderStripeHandle {
|
||||
DeferredReaderStripeHandle {
|
||||
state: Arc::clone(&self.state),
|
||||
stripe_stride,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Handle to a still-unopened [`DeferredObjectReader`] that can advance the
|
||||
/// pending open offset by whole bitrot blocks (stripes) before the first read.
|
||||
///
|
||||
/// The GET lockstep decode reads only the data shards while the object is
|
||||
/// healthy and keeps every parity slot as an unopened deferred reader. When a
|
||||
/// data shard dies at stripe `k`, the decoder uses this handle to shift the
|
||||
/// parity reader's pending open offset by `k` encoded blocks — the same
|
||||
/// `bitrot_encoded_range` geometry used when the reader was created — so its
|
||||
/// first read returns stripe `k` and the lockstep alignment invariant holds
|
||||
/// (backlog#923; alignment rule from backlog#832).
|
||||
///
|
||||
/// `stripe_stride` is `shard_size + checksum_algo.size()` per full stripe.
|
||||
/// For `hash_size == 0` (e.g. `HashAlgorithm::None`) this degrades to the
|
||||
/// identity mapping `k * shard_size`, matching `bitrot_encoded_range`.
|
||||
#[derive(Clone)]
|
||||
pub(crate) struct DeferredReaderStripeHandle {
|
||||
state: Arc<Mutex<DeferredObjectReaderState>>,
|
||||
stripe_stride: usize,
|
||||
}
|
||||
|
||||
impl DeferredReaderStripeHandle {
|
||||
/// Advance the pending source by `stripes` full stripes.
|
||||
///
|
||||
/// Returns `false` when the reader has already been opened (or failed):
|
||||
/// its stream position is then unknown to the caller and it must not be
|
||||
/// engaged mid-object.
|
||||
pub(crate) fn advance_stripes(&self, stripes: usize) -> bool {
|
||||
if stripes == 0 {
|
||||
return true;
|
||||
}
|
||||
let Ok(mut state) = self.state.lock() else {
|
||||
return false;
|
||||
};
|
||||
match &mut *state {
|
||||
DeferredObjectReaderState::Pending(Some(source)) => {
|
||||
let Some(delta) = self.stripe_stride.checked_mul(stripes) else {
|
||||
return false;
|
||||
};
|
||||
let Some(offset) = source.offset.checked_add(delta) else {
|
||||
return false;
|
||||
};
|
||||
source.offset = offset;
|
||||
source.length = source.length.saturating_sub(delta);
|
||||
true
|
||||
}
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -471,6 +528,38 @@ pub fn create_deferred_bitrot_reader(
|
||||
skip_verify: bool,
|
||||
use_mmap_read: bool,
|
||||
) -> BitrotReader<Box<dyn AsyncRead + Send + Sync + Unpin>> {
|
||||
create_deferred_bitrot_reader_with_stripe_handle(
|
||||
inline_data,
|
||||
disk,
|
||||
bucket,
|
||||
path,
|
||||
offset,
|
||||
length,
|
||||
shard_size,
|
||||
checksum_algo,
|
||||
skip_verify,
|
||||
use_mmap_read,
|
||||
)
|
||||
.0
|
||||
}
|
||||
|
||||
/// Like [`create_deferred_bitrot_reader`], but also returns a
|
||||
/// [`DeferredReaderStripeHandle`] that can realign the still-unopened reader
|
||||
/// to a later stripe before its first read (backlog#923).
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub(crate) fn create_deferred_bitrot_reader_with_stripe_handle(
|
||||
inline_data: Option<Bytes>,
|
||||
disk: Option<DiskStore>,
|
||||
bucket: &str,
|
||||
path: &str,
|
||||
offset: usize,
|
||||
length: usize,
|
||||
shard_size: usize,
|
||||
checksum_algo: HashAlgorithm,
|
||||
skip_verify: bool,
|
||||
use_mmap_read: bool,
|
||||
) -> (BitrotReader<Box<dyn AsyncRead + Send + Sync + Unpin>>, DeferredReaderStripeHandle) {
|
||||
let stripe_stride = shard_size + checksum_algo.size();
|
||||
let (offset, length) = bitrot_encoded_range(offset, length, shard_size, checksum_algo.clone());
|
||||
let source = BitrotReaderSource {
|
||||
inline_data,
|
||||
@@ -483,7 +572,15 @@ pub fn create_deferred_bitrot_reader(
|
||||
stage_metrics: None,
|
||||
};
|
||||
|
||||
BitrotReader::new(Box::new(DeferredObjectReader::new(source)), shard_size, checksum_algo, skip_verify)
|
||||
let deferred = DeferredObjectReader::new(source);
|
||||
let handle = deferred.stripe_handle(stripe_stride);
|
||||
let reader = BitrotReader::new(
|
||||
Box::new(deferred) as Box<dyn AsyncRead + Send + Sync + Unpin>,
|
||||
shard_size,
|
||||
checksum_algo,
|
||||
skip_verify,
|
||||
);
|
||||
(reader, handle)
|
||||
}
|
||||
|
||||
/// Create a new BitrotWriterWrapper based on the provided parameters
|
||||
@@ -741,6 +838,139 @@ mod tests {
|
||||
assert_eq!(&out[..n], b"efgh");
|
||||
}
|
||||
|
||||
/// Merge gate for backlog#923: engaging a parity shard at stripe `k`
|
||||
/// converts `k` stripes into an encoded byte offset via the
|
||||
/// `bitrot_encoded_range` geometry. Cover both on-disk formats: the
|
||||
/// streaming checksum layout (32-byte hash per block) and the
|
||||
/// no-per-block-hash layout (`hash_size == 0`), where the mapping must
|
||||
/// degrade to the identity `k * shard_size`.
|
||||
#[tokio::test]
|
||||
async fn test_deferred_stripe_handle_aligns_reader_to_requested_stripe() {
|
||||
let shard_size = 4;
|
||||
let payload = b"aaaabbbbccccdddd"; // 4 full bitrot blocks
|
||||
|
||||
for algo in [HashAlgorithm::HighwayHash256S, HashAlgorithm::None] {
|
||||
let mut writer =
|
||||
create_bitrot_writer(true, None, "test-volume", "test-path", payload.len() as i64, shard_size, algo.clone())
|
||||
.await
|
||||
.expect("inline bitrot writer");
|
||||
for chunk in payload.chunks(shard_size) {
|
||||
writer.write(chunk).await.expect("write chunk");
|
||||
}
|
||||
let inline_data = writer.into_inline_data().expect("inline buffer");
|
||||
|
||||
for stripe in 0..payload.len() / shard_size {
|
||||
let (mut reader, handle) = create_deferred_bitrot_reader_with_stripe_handle(
|
||||
Some(inline_data.clone().into()),
|
||||
None,
|
||||
"test-bucket",
|
||||
"test-path",
|
||||
0,
|
||||
payload.len(),
|
||||
shard_size,
|
||||
algo.clone(),
|
||||
false,
|
||||
false,
|
||||
);
|
||||
|
||||
assert!(
|
||||
handle.advance_stripes(stripe),
|
||||
"pending deferred reader must accept a stripe advance (algo={algo:?}, stripe={stripe})"
|
||||
);
|
||||
|
||||
let mut out = [0u8; 4];
|
||||
let n = reader.read(&mut out).await.expect("read shard block after stripe advance");
|
||||
assert_eq!(n, shard_size);
|
||||
assert_eq!(
|
||||
&out[..n],
|
||||
&payload[stripe * shard_size..(stripe + 1) * shard_size],
|
||||
"advanced reader must return exactly stripe {stripe} (algo={algo:?})"
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Bitrot verification must stay active for a parity shard engaged
|
||||
/// mid-object: after `advance_stripes(k)` the reader verifies stripe `k`'s
|
||||
/// block against stripe `k`'s stored hash.
|
||||
#[tokio::test]
|
||||
async fn test_deferred_stripe_handle_preserves_bitrot_verification_after_advance() {
|
||||
let shard_size = 4;
|
||||
let algo = HashAlgorithm::HighwayHash256S;
|
||||
let payload = b"aaaabbbbccccdddd";
|
||||
|
||||
let mut writer =
|
||||
create_bitrot_writer(true, None, "test-volume", "test-path", payload.len() as i64, shard_size, algo.clone())
|
||||
.await
|
||||
.expect("inline bitrot writer");
|
||||
for chunk in payload.chunks(shard_size) {
|
||||
writer.write(chunk).await.expect("write chunk");
|
||||
}
|
||||
let mut inline_data = writer.into_inline_data().expect("inline buffer");
|
||||
|
||||
// Corrupt one payload byte of block 2 (blocks are hash + data).
|
||||
let block_len = algo.size() + shard_size;
|
||||
inline_data[2 * block_len + algo.size()] ^= 0x01;
|
||||
|
||||
let (mut reader, handle) = create_deferred_bitrot_reader_with_stripe_handle(
|
||||
Some(inline_data.into()),
|
||||
None,
|
||||
"test-bucket",
|
||||
"test-path",
|
||||
0,
|
||||
payload.len(),
|
||||
shard_size,
|
||||
algo,
|
||||
false,
|
||||
false,
|
||||
);
|
||||
assert!(handle.advance_stripes(2));
|
||||
|
||||
let mut out = [0u8; 4];
|
||||
let err = reader
|
||||
.read(&mut out)
|
||||
.await
|
||||
.expect_err("bitrot mismatch at the advanced stripe must fail the read");
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
}
|
||||
|
||||
/// A reader that has already been opened has an unknown stream position
|
||||
/// from the handle's point of view; the advance must be refused so the
|
||||
/// caller retires the reader instead of engaging it out of alignment.
|
||||
#[tokio::test]
|
||||
async fn test_deferred_stripe_handle_rejects_advance_after_open() {
|
||||
let shard_size = 4;
|
||||
let algo = HashAlgorithm::HighwayHash256S;
|
||||
let payload = b"aaaabbbb";
|
||||
|
||||
let mut writer =
|
||||
create_bitrot_writer(true, None, "test-volume", "test-path", payload.len() as i64, shard_size, algo.clone())
|
||||
.await
|
||||
.expect("inline bitrot writer");
|
||||
for chunk in payload.chunks(shard_size) {
|
||||
writer.write(chunk).await.expect("write chunk");
|
||||
}
|
||||
let inline_data = writer.into_inline_data().expect("inline buffer");
|
||||
|
||||
let (mut reader, handle) = create_deferred_bitrot_reader_with_stripe_handle(
|
||||
Some(inline_data.into()),
|
||||
None,
|
||||
"test-bucket",
|
||||
"test-path",
|
||||
0,
|
||||
payload.len(),
|
||||
shard_size,
|
||||
algo,
|
||||
false,
|
||||
false,
|
||||
);
|
||||
|
||||
let mut out = [0u8; 4];
|
||||
reader.read(&mut out).await.expect("first read opens the deferred source");
|
||||
|
||||
assert!(!handle.advance_stripes(1), "an opened deferred reader must reject stripe advances");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_create_bitrot_reader_without_data_or_disk() {
|
||||
let shard_size = 16;
|
||||
|
||||
@@ -48,7 +48,8 @@ use crate::diagnostics::get::{
|
||||
};
|
||||
use crate::erasure::coding::BitrotReader;
|
||||
use crate::io_support::bitrot::{
|
||||
BitrotReaderStageMetrics, create_bitrot_reader_with_stage_metrics, create_deferred_bitrot_reader, object_mmap_read_enabled,
|
||||
BitrotReaderStageMetrics, DeferredReaderStripeHandle, create_bitrot_reader_with_stage_metrics,
|
||||
create_deferred_bitrot_reader_with_stripe_handle, object_mmap_read_enabled,
|
||||
};
|
||||
use crate::set_disk::shard_source::ShardReadCost;
|
||||
use futures::stream::{FuturesUnordered, StreamExt};
|
||||
@@ -911,6 +912,10 @@ pub(in crate::set_disk) const DIRECT_MEMORY_BITROT_READER_STAGE_METRICS: BitrotR
|
||||
|
||||
pub(in crate::set_disk) struct BitrotReaderSetup {
|
||||
pub(in crate::set_disk) readers: Vec<Option<ObjectBitrotReader>>,
|
||||
/// Per-slot stripe handles for readers that are still unopened deferred
|
||||
/// readers. The lockstep GET decode uses them to open a parity shard
|
||||
/// aligned to the stripe where a data shard failed (backlog#923).
|
||||
pub(in crate::set_disk) deferred_stripe_handles: Vec<Option<DeferredReaderStripeHandle>>,
|
||||
pub(in crate::set_disk) errors: Vec<Option<DiskError>>,
|
||||
pub(in crate::set_disk) scheduled: Vec<bool>,
|
||||
pub(in crate::set_disk) attempted: Vec<bool>,
|
||||
@@ -985,6 +990,7 @@ impl BitrotReaderSetup {
|
||||
pub(in crate::set_disk) fn new(shards: usize) -> Self {
|
||||
Self {
|
||||
readers: (0..shards).map(|_| None).collect(),
|
||||
deferred_stripe_handles: (0..shards).map(|_| None).collect(),
|
||||
errors: vec![Some(DiskError::DiskNotFound); shards],
|
||||
scheduled: vec![false; shards],
|
||||
attempted: vec![false; shards],
|
||||
@@ -1110,8 +1116,14 @@ impl BitrotReaderSetup {
|
||||
}
|
||||
}
|
||||
|
||||
pub(in crate::set_disk) fn retain_deferred_reader(&mut self, idx: usize, reader: ObjectBitrotReader) {
|
||||
pub(in crate::set_disk) fn retain_deferred_reader(
|
||||
&mut self,
|
||||
idx: usize,
|
||||
reader: ObjectBitrotReader,
|
||||
stripe_handle: DeferredReaderStripeHandle,
|
||||
) {
|
||||
self.readers[idx] = Some(reader);
|
||||
self.deferred_stripe_handles[idx] = Some(stripe_handle);
|
||||
self.errors[idx] = None;
|
||||
self.deferred_count = self.deferred_count.saturating_add(1);
|
||||
}
|
||||
@@ -1204,21 +1216,58 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
|
||||
let disk = disks[idx].clone();
|
||||
let data_dir = files[idx].data_dir.unwrap_or_default();
|
||||
let path = format!("{object}/{data_dir}/part.{part_number}");
|
||||
setup.retain_deferred_reader(
|
||||
idx,
|
||||
create_deferred_bitrot_reader(
|
||||
inline_data,
|
||||
disk,
|
||||
bucket,
|
||||
&path,
|
||||
read_offset,
|
||||
read_length,
|
||||
shard_size,
|
||||
checksum_algo.clone(),
|
||||
skip_verify_bitrot,
|
||||
use_mmap_read,
|
||||
),
|
||||
let (reader, stripe_handle) = create_deferred_bitrot_reader_with_stripe_handle(
|
||||
inline_data,
|
||||
disk,
|
||||
bucket,
|
||||
&path,
|
||||
read_offset,
|
||||
read_length,
|
||||
shard_size,
|
||||
checksum_algo.clone(),
|
||||
skip_verify_bitrot,
|
||||
use_mmap_read,
|
||||
);
|
||||
setup.retain_deferred_reader(idx, reader, stripe_handle);
|
||||
}
|
||||
|
||||
// With the data-shards-only lockstep gate on (backlog#923), the GET decode
|
||||
// reads only the data shards while the object is healthy; a parity reader
|
||||
// is engaged on demand and must therefore stay unopened so its start
|
||||
// offset can be advanced to the failing stripe. A parity reader that was
|
||||
// opened eagerly during setup is pinned at the stripe-0 stream position
|
||||
// and could never be engaged mid-object, so swap it for an unopened
|
||||
// deferred reader carrying a stripe handle. The eager open already proved
|
||||
// the shard is reachable; no shard bytes were read from it, and the
|
||||
// ready/error bookkeeping that quorum decisions rely on is left untouched.
|
||||
// Gate off (default): keep the eagerly opened parity readers exactly as
|
||||
// before — the lockstep path reads them on every stripe.
|
||||
if !crate::erasure::coding::decode::get_lockstep_data_shards_only_enabled() {
|
||||
return;
|
||||
}
|
||||
for idx in data_shards..disks.len() {
|
||||
if setup.readers[idx].is_none() || setup.deferred_stripe_handles[idx].is_some() {
|
||||
continue;
|
||||
}
|
||||
|
||||
let inline_data = files[idx].data.clone();
|
||||
let disk = disks[idx].clone();
|
||||
let data_dir = files[idx].data_dir.unwrap_or_default();
|
||||
let path = format!("{object}/{data_dir}/part.{part_number}");
|
||||
let (reader, stripe_handle) = create_deferred_bitrot_reader_with_stripe_handle(
|
||||
inline_data,
|
||||
disk,
|
||||
bucket,
|
||||
&path,
|
||||
read_offset,
|
||||
read_length,
|
||||
shard_size,
|
||||
checksum_algo.clone(),
|
||||
skip_verify_bitrot,
|
||||
use_mmap_read,
|
||||
);
|
||||
setup.readers[idx] = Some(reader);
|
||||
setup.deferred_stripe_handles[idx] = Some(stripe_handle);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -912,13 +912,18 @@ impl SetDisks {
|
||||
let decode_stage_start = Instant::now();
|
||||
let unattempted_data_shards = !reader_setup.data_shards_attempted(erasure.data_shards);
|
||||
let readers = reader_setup.readers;
|
||||
let (written, err) = if let Some(read_costs) = read_costs {
|
||||
erasure
|
||||
.decode_with_read_costs(writer, readers, part_offset, part_length, part_size, read_costs)
|
||||
.await
|
||||
} else {
|
||||
erasure.decode(writer, readers, part_offset, part_length, part_size).await
|
||||
};
|
||||
let deferred_stripe_handles = reader_setup.deferred_stripe_handles;
|
||||
let (written, err) = erasure
|
||||
.decode_with_stripe_handles(
|
||||
writer,
|
||||
readers,
|
||||
part_offset,
|
||||
part_length,
|
||||
part_size,
|
||||
read_costs,
|
||||
deferred_stripe_handles,
|
||||
)
|
||||
.await;
|
||||
let decode_elapsed = decode_stage_start.elapsed();
|
||||
rustfs_io_metrics::record_get_object_decode_duration(decode_elapsed.as_secs_f64());
|
||||
rustfs_io_metrics::record_get_object_stage_duration_by_size(
|
||||
@@ -1271,6 +1276,7 @@ impl SetDisks {
|
||||
}
|
||||
|
||||
let readers = reader_setup.readers;
|
||||
let deferred_stripe_handles = reader_setup.deferred_stripe_handles;
|
||||
let source = if let Some(read_costs) = read_costs {
|
||||
coding::decode::ParallelReader::new_with_metrics_path_read_costs_and_reconstruction_verification(
|
||||
readers,
|
||||
@@ -1288,7 +1294,8 @@ impl SetDisks {
|
||||
part_size,
|
||||
Some(metrics_path),
|
||||
)
|
||||
};
|
||||
}
|
||||
.with_deferred_parity_handles(deferred_stripe_handles);
|
||||
let engine = build_get_codec_streaming_decode_engine(erasure.clone())?;
|
||||
let reader =
|
||||
coding::decode_reader::ErasureDecodeReader::new_with_metrics_path(source, engine, part_length, metrics_path)?;
|
||||
@@ -3205,6 +3212,60 @@ mod tests {
|
||||
assert_eq!(&out[..n], [b"aaaa", b"bbbb", b"cccc", b"dddd"][fallback_index]);
|
||||
}
|
||||
|
||||
/// backlog#923: with the data-shards-only lockstep gate on, every retained
|
||||
/// 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
|
||||
/// the gate off (default), eagerly opened parity readers are kept exactly
|
||||
/// as before and carry no handles.
|
||||
#[tokio::test]
|
||||
#[serial_test::serial]
|
||||
async fn bitrot_reader_setup_gates_parity_stripe_handle_conversion() {
|
||||
for enabled in [None, Some("true")] {
|
||||
// A missing data shard forces the VerifyReconstruction quorum to 3,
|
||||
// so both parity slots complete eagerly (attempted + ready) before
|
||||
// the deferred fill runs.
|
||||
let mut setup = temp_env::async_with_vars(
|
||||
[("RUSTFS_GET_LOCKSTEP_DATA_SHARDS_ONLY_ENABLE", enabled)],
|
||||
setup_inline_bitrot_readers_with_preference(
|
||||
vec![None, Some(b"bbbb"), Some(b"cccc"), Some(b"dddd")],
|
||||
2,
|
||||
2,
|
||||
BitrotReaderSetupMode::VerifyReconstruction,
|
||||
false,
|
||||
),
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(setup.available_shards(), 3);
|
||||
for idx in 2..4 {
|
||||
assert!(setup.attempted[idx] && setup.ready[idx], "parity slot {idx} should be eagerly ready");
|
||||
assert!(setup.readers[idx].is_some(), "parity slot {idx} must keep a reader (enabled={enabled:?})");
|
||||
assert_eq!(
|
||||
setup.deferred_stripe_handles[idx].is_some(),
|
||||
enabled.is_some(),
|
||||
"parity slot {idx} stripe handle must match the gate (enabled={enabled:?})"
|
||||
);
|
||||
}
|
||||
|
||||
if enabled.is_some() {
|
||||
// The converted parity reader is still unopened: its handle
|
||||
// accepts a stripe advance, and reading it yields the shard
|
||||
// bytes (block 0 here).
|
||||
let handle = setup.deferred_stripe_handles[3].as_ref().expect("slot 3 handle");
|
||||
assert!(handle.advance_stripes(1), "unopened converted parity reader must accept a stripe advance");
|
||||
|
||||
let mut reader = setup.readers[2].take().expect("slot 2 reader");
|
||||
let mut out = [0u8; 4];
|
||||
let n = reader
|
||||
.read(&mut out)
|
||||
.await
|
||||
.expect("converted parity reader should open on read");
|
||||
assert_eq!(n, 4);
|
||||
assert_eq!(&out[..n], b"cccc");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn bitrot_reader_setup_data_blocks_first_keeps_deferred_fallback_readers() {
|
||||
let mut setup = setup_inline_bitrot_readers_with_env(
|
||||
|
||||
Reference in New Issue
Block a user