fix: bound restarted-peer stalls on quorum reads and writes (#8152)

A restarted peer can accept a pooled connection and never send response
headers, so HttpReader::open waited past the client body timeout before
the body-stall timer or erasure hedge could run. Bound that header wait
by the stall timeout and retry the open once on a fresh connection.

After write quorum, MultiWriter still waited out the full disk stall
for a silent peer, which matches the client timeout. Give remaining
writers one second, then drop them so the caller returns.

Signed-off-by: loverustfs <155562731+loverustfs@users.noreply.github.com>
This commit is contained in:
RustFS
2026-09-28 10:22:41 +08:00
committed by GitHub
parent 719000460b
commit 2e014d25b3
4 changed files with 341 additions and 15 deletions
+46 -3
View File
@@ -993,7 +993,19 @@ impl RemoteDisk {
// internode transport failure classified as retryable can be safely re-dialed. The
// classifier is direction-agnostic — it inspects the InternodeHttpError kind — so it is
// reused here from the write path.
err.is_retryable_internode_write_failure()
if err.is_retryable_internode_write_failure() {
return true;
}
// Header wait uses the body stall budget. That error is `BodyStalled`, not an
// `InternodeHttpError`, and the open has not consumed a shard byte yet.
matches!(
err,
DiskError::Io(error)
if error
.get_ref()
.and_then(|source| source.downcast_ref::<rustfs_rio::BodyStalled>())
.is_some()
)
}
pub(crate) async fn new(ep: &Endpoint, opt: &DiskOption, data_transport: Arc<dyn InternodeDataTransport>) -> Result<Self> {
@@ -1102,7 +1114,14 @@ impl RemoteDisk {
let mut attempt = 1;
let mut last_retry_classification = None;
loop {
match self.data_transport.open_read(request.clone()).await {
// The second attempt bypasses the pool. The first failure is often a
// stale kept-alive connection to a peer that just restarted.
let opened = if attempt == 1 {
self.data_transport.open_read(request.clone()).await
} else {
self.data_transport.open_read_fresh(request.clone()).await
};
match opened {
Ok(reader) => {
if attempt > 1
&& let Some(classification) = last_retry_classification
@@ -1138,7 +1157,12 @@ impl RemoteDisk {
let mut attempt = 1;
let mut last_retry_classification = None;
loop {
match self.data_transport.open_read_chunks(request.clone()).await {
let opened = if attempt == 1 {
self.data_transport.open_read_chunks(request.clone()).await
} else {
self.data_transport.open_read_chunks_fresh(request.clone()).await
};
match opened {
Ok(reader) => {
if attempt > 1
&& let Some(classification) = last_retry_classification
@@ -7729,6 +7753,25 @@ mod tests {
assert_eq!(transport.calls().len(), 2, "read_file_stream should retry exactly once");
}
#[tokio::test]
async fn test_remote_disk_read_file_stream_retries_header_stall_on_fresh_connection() {
let stall = rustfs_rio::BodyStalled {
timeout: Duration::from_millis(50),
};
let transport = RetryingOpenReadInternodeDataTransport::with_steps(vec![
OpenWriteTestStep::Error(DiskError::Io(std::io::Error::new(std::io::ErrorKind::TimedOut, stall))),
OpenWriteTestStep::Success,
]);
let remote_disk = new_remote_disk_with_transport(Arc::new(transport.clone())).await;
let _reader = remote_disk
.read_file_stream("bucket", "object/part.1", 0, 4096)
.await
.expect("a header stall must be retried once before the shard is failed");
assert_eq!(transport.calls().len(), 2, "header stall should re-dial exactly once");
}
#[tokio::test]
async fn test_remote_disk_read_file_stream_does_not_retry_non_retryable_open_read_error() {
let transport = RetryingOpenReadInternodeDataTransport::with_steps(vec![OpenWriteTestStep::Error(DiskError::from(
@@ -5267,6 +5267,71 @@ mod tests {
assert_eq!(DATA_SHARDS + 1, bufs.iter().filter(|buf| buf.is_some()).count());
}
/// A peer that sends part of a shard and then stops must not pin the stripe
/// for the full read timeout. The other EC 2+2 shards already hold a
/// decode-plus-verification quorum, so the lockstep hedge retires the
/// stalled reader and the stripe completes.
#[tokio::test]
async fn test_lockstep_hedges_shard_that_stops_mid_read() {
const NUM_SHARDS: usize = 1;
const BLOCK_SIZE: usize = 64;
const DATA_SHARDS: usize = 2;
const PARITY_SHARDS: usize = 2;
const SHARD_SIZE: usize = BLOCK_SIZE / DATA_SHARDS;
let hash_algo = HashAlgorithm::None;
let readers = vec![
Some(BitrotReader::new(
TestShardReader::PartialThenPending {
data: vec![0xab],
emitted: false,
},
SHARD_SIZE,
hash_algo.clone(),
false,
)),
Some(BitrotReader::new(
TestShardReader::Ready(Cursor::new(vec![1_u8; SHARD_SIZE * NUM_SHARDS])),
SHARD_SIZE,
hash_algo.clone(),
false,
)),
Some(BitrotReader::new(
TestShardReader::Ready(Cursor::new(vec![2_u8; SHARD_SIZE * NUM_SHARDS])),
SHARD_SIZE,
hash_algo.clone(),
false,
)),
Some(BitrotReader::new(
TestShardReader::Ready(Cursor::new(vec![3_u8; SHARD_SIZE * NUM_SHARDS])),
SHARD_SIZE,
hash_algo,
false,
)),
];
let erasure = Erasure::new(DATA_SHARDS, PARITY_SHARDS, BLOCK_SIZE);
let mut parallel_reader = ParallelReader::new_with_metrics_path_read_costs_timeout_and_reconstruction_verification(
readers,
erasure,
0,
NUM_SHARDS * BLOCK_SIZE,
None,
vec![ShardReadCost::Unknown; 4],
Duration::from_secs(60),
true,
);
let (bufs, errs) = tokio::time::timeout(Duration::from_millis(500), parallel_reader.read())
.await
.expect("a shard that stops mid-read must be hedged instead of waiting out the 60s read timeout");
assert!(matches!(&errs[0], Some(DiskError::Io(err)) if err.kind() == ErrorKind::TimedOut));
assert!(parallel_reader.readers[0].is_none());
assert!(bufs[0].is_none());
assert_eq!(DATA_SHARDS + 1, bufs.iter().filter(|buf| buf.is_some()).count());
}
#[tokio::test(start_paused = true)]
async fn lockstep_reopens_hedged_shard_after_later_peer_loss() {
const DATA_SHARDS: usize = 2;
+164 -11
View File
@@ -25,6 +25,7 @@ use bytes::{Bytes, BytesMut};
use futures::StreamExt;
use futures::stream::FuturesUnordered;
use rustfs_utils::HashAlgorithm;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Instant;
use std::vec;
@@ -325,6 +326,80 @@ impl Default for WriteProgressPolicy {
}
}
/// After write quorum is already in hand, how long a slower shard may still
/// finish before it is dropped. The per-shard stall budget (30s by default)
/// is what a black-hole peer is allowed to consume *before* quorum exists.
/// Once quorum exists, waiting out that whole budget pins the caller for as
/// long as the client request timeout, so the client gives up on a write the
/// remaining disks have already accepted.
const WRITE_QUORUM_STRAGGLER_GRACE: std::time::Duration = std::time::Duration::from_secs(1);
fn post_quorum_straggler_grace(stall_timeout: Option<std::time::Duration>) -> Option<std::time::Duration> {
stall_timeout.map(|stall| stall.min(WRITE_QUORUM_STRAGGLER_GRACE))
}
/// Drive shard ops until write quorum has succeeded and stragglers have had
/// [`post_quorum_straggler_grace`] to finish. Returns `true` for each slot
/// whose future completed; `false` means the caller must fail that slot.
async fn await_write_quorum<F>(
write_quorum: usize,
stall_timeout: Option<std::time::Duration>,
slot_count: usize,
mut futures: FuturesUnordered<F>,
) -> Vec<bool>
where
F: std::future::Future<Output = (usize, bool)>,
{
let mut saw = vec![false; slot_count];
let mut successes = 0usize;
let mut grace: Option<Pin<Box<tokio::time::Sleep>>> = None;
loop {
if futures.is_empty() {
break;
}
let completed = match grace.as_mut() {
Some(sleep) => {
tokio::select! {
biased;
item = futures.next() => item,
() = sleep.as_mut() => break,
}
}
None => futures.next().await,
};
let Some((index, succeeded)) = completed else {
break;
};
if let Some(slot) = saw.get_mut(index) {
*slot = true;
}
if succeeded {
successes = successes.saturating_add(1);
if successes >= write_quorum
&& grace.is_none()
&& let Some(delay) = post_quorum_straggler_grace(stall_timeout)
{
grace = Some(Box::pin(tokio::time::sleep(delay)));
}
}
}
saw
}
fn abandon_unfinished_writers(writers: &mut [Option<BitrotWriterWrapper>], errs: &mut [Option<Error>], saw: &[bool]) {
for (index, seen) in saw.iter().enumerate() {
if *seen || errs.get(index).and_then(Option::as_ref).is_some() {
continue;
}
if let Some(err) = errs.get_mut(index) {
*err = Some(Error::Timeout);
}
if let Some(writer) = writers.get_mut(index) {
*writer = None;
}
}
}
pub(crate) struct MultiWriter<'a> {
writers: &'a mut [Option<BitrotWriterWrapper>],
integrity: Option<&'a mut IntegrityBuilder>,
@@ -418,9 +493,12 @@ impl<'a> MultiWriter<'a> {
assert_eq!(shards.len(), self.writers.len());
let budget = self.next_progress_budget();
{
let mut futures = FuturesUnordered::new();
for ((writer_opt, err), shard) in self.writers.iter_mut().zip(self.errs.iter_mut()).zip(shards) {
let stall_timeout = self.policy.stall_timeout;
let write_quorum = self.write_quorum;
let slot_count = self.writers.len();
let saw = {
let futures = FuturesUnordered::new();
for (index, ((writer_opt, err), shard)) in self.writers.iter_mut().zip(self.errs.iter_mut()).zip(shards).enumerate() {
if err.is_some() {
continue; // Skip if we already have an error for this writer
}
@@ -428,7 +506,9 @@ impl<'a> MultiWriter<'a> {
// failed and its disk dropped, so a stalled peer cannot pin an
// otherwise-healthy write quorum (rustfs/backlog#1319). `budget`
// is recomputed per block, so it bounds a stall — not the total
// transfer time — and a slow-but-honest writer is never killed.
// transfer time — and a slow-but-honest writer is never killed
// while quorum is still open. Once quorum has succeeded, a
// straggler only gets `WRITE_QUORUM_STRAGGLER_GRACE` more time.
futures.push(async move {
match budget {
Some(budget) => match tokio::time::timeout(budget, Self::write_shard(writer_opt, err, shard)).await {
@@ -440,10 +520,13 @@ impl<'a> MultiWriter<'a> {
},
None => Self::write_shard(writer_opt, err, shard).await,
}
let succeeded = err.is_none() && writer_opt.is_some();
(index, succeeded)
});
}
while let Some(()) = futures.next().await {}
}
await_write_quorum(write_quorum, stall_timeout, slot_count, futures).await
};
abandon_unfinished_writers(self.writers, &mut self.errs, &saw);
let nil_count = self.errs.iter().filter(|&e| e.is_none()).count();
if nil_count >= self.write_quorum {
@@ -490,9 +573,12 @@ impl<'a> MultiWriter<'a> {
pub async fn shutdown(&mut self) -> std::io::Result<()> {
crate::hp_guard!("MultiWriter::shutdown");
let budget = self.next_progress_budget();
{
let mut futures = FuturesUnordered::new();
for (writer_opt, err) in self.writers.iter_mut().zip(self.errs.iter_mut()) {
let stall_timeout = self.policy.stall_timeout;
let write_quorum = self.write_quorum;
let slot_count = self.writers.len();
let saw = {
let futures = FuturesUnordered::new();
for (index, (writer_opt, err)) in self.writers.iter_mut().zip(self.errs.iter_mut()).enumerate() {
if err.is_some() {
continue;
}
@@ -501,6 +587,8 @@ impl<'a> MultiWriter<'a> {
// here forever for a small object whose bytes were fully buffered
// (so `write` never blocked). Bound it with the same progress
// budget and drop the stalled writer before the quorum check.
// After quorum has shut down, do not keep waiting out the rest of
// that budget for a peer that has stopped answering.
futures.push(async move {
match budget {
Some(budget) => match tokio::time::timeout(budget, Self::shutdown_writer(writer_opt, err)).await {
@@ -512,10 +600,13 @@ impl<'a> MultiWriter<'a> {
},
None => Self::shutdown_writer(writer_opt, err).await,
}
let succeeded = err.is_none() && writer_opt.is_some();
(index, succeeded)
});
}
while let Some(()) = futures.next().await {}
}
await_write_quorum(write_quorum, stall_timeout, slot_count, futures).await
};
abandon_unfinished_writers(self.writers, &mut self.errs, &saw);
let nil_count = self.errs.iter().filter(|&e| e.is_none()).count();
if nil_count >= self.write_quorum {
@@ -2134,6 +2225,68 @@ mod tests {
assert!(writers[1].is_some() && writers[2].is_some() && writers[3].is_some());
}
/// Quorum is met by the three healthy writers at t=0. The black-hole peer
/// must be abandoned after the post-quorum grace, not after the full 5s
/// stall budget. Waiting out the stall makes a 30s default stall consume
/// the whole client request budget.
#[tokio::test(start_paused = true)]
async fn multi_writer_quorum_abandons_black_hole_after_straggler_grace() {
let committed: Vec<Arc<Mutex<Vec<u8>>>> = (0..3).map(|_| Arc::new(Mutex::new(Vec::new()))).collect();
let mut writers = vec![
Some(bitrot_writer_plain(StallingWriter::stalls_on_write(0), 64)),
Some(bitrot_writer_plain(DeferredCommitWriter::new(committed[0].clone()), 64)),
Some(bitrot_writer_plain(DeferredCommitWriter::new(committed[1].clone()), 64)),
Some(bitrot_writer_plain(DeferredCommitWriter::new(committed[2].clone()), 64)),
];
let started = tokio::time::Instant::now();
{
let policy = WriteProgressPolicy::new(Duration::from_secs(5), Duration::ZERO);
let mut mw = MultiWriter::with_policy(&mut writers, 3, policy);
mw.write(four_shards())
.await
.expect("quorum must succeed without waiting out the black-hole stall");
}
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_secs(2),
"post-quorum grace must cut the black-hole wait short of the 5s stall, elapsed {elapsed:?}"
);
assert!(
elapsed >= WRITE_QUORUM_STRAGGLER_GRACE,
"the healthy quorum must still give the straggler its grace window, elapsed {elapsed:?}"
);
assert!(writers[0].is_none(), "the black-hole writer must be dropped once grace expires");
assert!(writers[1].is_some() && writers[2].is_some() && writers[3].is_some());
}
/// A slow-but-honest shard that finishes inside the post-quorum grace keeps
/// its writer. The grace exists to drop peers that have stopped, not to
/// discard a disk that is merely behind the fastest quorum.
#[tokio::test(start_paused = true)]
async fn multi_writer_quorum_keeps_straggler_that_finishes_within_grace() {
let committed: Vec<Arc<Mutex<Vec<u8>>>> = (0..4).map(|_| Arc::new(Mutex::new(Vec::new()))).collect();
let mut writers = vec![
Some(bitrot_writer_plain(SlowWriter::new(Duration::from_millis(200)), 64)),
Some(bitrot_writer_plain(DeferredCommitWriter::new(committed[1].clone()), 64)),
Some(bitrot_writer_plain(DeferredCommitWriter::new(committed[2].clone()), 64)),
Some(bitrot_writer_plain(DeferredCommitWriter::new(committed[3].clone()), 64)),
];
{
let policy = WriteProgressPolicy::new(Duration::from_secs(5), Duration::ZERO);
let mut mw = MultiWriter::with_policy(&mut writers, 3, policy);
mw.write(four_shards())
.await
.expect("a straggler inside the grace window must still satisfy the write");
}
assert!(
writers.iter().all(Option::is_some),
"a shard that finishes within the grace window must not be dropped"
);
}
// Two black-hole writers drop the healthy count to 2/4, below the quorum of
// 3, so the write must fail cleanly (not hang).
#[tokio::test(start_paused = true)]
+66 -1
View File
@@ -1065,7 +1065,26 @@ impl HttpReader {
}
let request_started = Instant::now();
let resp = request.send().await.map_err(|e| {
// `send()` resolves at the response headers, before the body stall timer
// in `poll_read` can run. A restarted peer (or a pooled connection left
// half-open when its pod network namespace disappeared) accepts the TCP
// connection and then never sends headers. Without this bound the shard
// open waits out kernel retransmits, long after the client has given up
// on the GET. The body stall budget is the same deadline: a header
// black hole is the same failure as a body that stops mid-shard.
let send_result = match stall_timeout {
Some(stall_timeout) => match time::timeout(stall_timeout, request.send()).await {
Ok(result) => result,
Err(_elapsed) => {
record_internode_operation_duration(track_internode_metrics, internode_operation, request_started.elapsed());
record_internode_stall_timeout(track_internode_metrics, internode_operation);
record_internode_error(track_internode_metrics, internode_operation);
return Err(body_stalled_error(stall_timeout));
}
},
None => request.send().await,
};
let resp = send_result.map_err(|e| {
record_internode_operation_duration(track_internode_metrics, internode_operation, request_started.elapsed());
record_internode_error(track_internode_metrics, internode_operation);
record_internode_classified_error(track_internode_metrics, internode_operation, classify_reqwest_error(&e));
@@ -2672,6 +2691,52 @@ mod tests {
handle.abort();
}
/// A peer that accepts the connection and then never sends response headers.
/// This is the restarted-pod case: the pooled TCP connection stays open, so
/// connect timeout does not fire, and the body stall timer has not started
/// because `send()` has not returned. The open itself must fail as
/// `BodyStalled` inside the stall budget.
#[tokio::test]
async fn http_reader_header_stall_fails_open_within_stall_budget() {
let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
Ok(listener) => listener,
Err(err) if err.kind() == std::io::ErrorKind::PermissionDenied => return,
Err(err) => panic!("test listener should bind: {err}"),
};
let addr = listener.local_addr().expect("listener local address should be available");
let app = Router::new().route(
"/hang-headers",
axum::routing::get(|| async {
std::future::pending::<()>().await;
StatusCode::OK
}),
);
let server_handle = tokio::spawn(async move {
axum::serve(listener, app).await.unwrap();
});
let url = format!("http://{addr}/hang-headers");
let stall = Duration::from_millis(50);
let opened = tokio::time::timeout(
Duration::from_secs(2),
HttpReader::new_with_stall_timeout(url, Method::GET, HeaderMap::new(), None, Some(stall)),
)
.await
.expect("header stall must fail the open instead of hanging until the test deadline");
let err = match opened {
Ok(_reader) => panic!("a peer that never sends headers must fail the reader open"),
Err(err) => err,
};
assert_eq!(err.kind(), io::ErrorKind::TimedOut);
let stalled = err
.get_ref()
.and_then(|source| source.downcast_ref::<BodyStalled>())
.expect("header stall should retain the typed body-stalled source");
assert_eq!(stalled.timeout, stall);
server_handle.abort();
}
#[tokio::test]
async fn http_chunk_reader_stall_timeout_retains_typed_source() {
let state = TestState::default();