From 2e014d25b3622418eefc17617afba293f3fac8db Mon Sep 17 00:00:00 2001 From: RustFS Date: Mon, 28 Sep 2026 10:22:41 +0800 Subject: [PATCH] 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> --- crates/ecstore/src/cluster/rpc/remote_disk.rs | 49 ++++- crates/ecstore/src/erasure/coding/decode.rs | 65 +++++++ crates/ecstore/src/erasure/coding/encode.rs | 175 ++++++++++++++++-- crates/rio/src/http_reader.rs | 67 ++++++- 4 files changed, 341 insertions(+), 15 deletions(-) diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index a4cfe9945..1f8cb9042 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -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::()) + .is_some() + ) } pub(crate) async fn new(ep: &Endpoint, opt: &DiskOption, data_transport: Arc) -> Result { @@ -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( diff --git a/crates/ecstore/src/erasure/coding/decode.rs b/crates/ecstore/src/erasure/coding/decode.rs index ab3776b6c..ad9433218 100644 --- a/crates/ecstore/src/erasure/coding/decode.rs +++ b/crates/ecstore/src/erasure/coding/decode.rs @@ -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; diff --git a/crates/ecstore/src/erasure/coding/encode.rs b/crates/ecstore/src/erasure/coding/encode.rs index 808105cfe..946e665f0 100644 --- a/crates/ecstore/src/erasure/coding/encode.rs +++ b/crates/ecstore/src/erasure/coding/encode.rs @@ -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) -> Option { + 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( + write_quorum: usize, + stall_timeout: Option, + slot_count: usize, + mut futures: FuturesUnordered, +) -> Vec +where + F: std::future::Future, +{ + let mut saw = vec![false; slot_count]; + let mut successes = 0usize; + let mut grace: Option>> = 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], errs: &mut [Option], 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], 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>>> = (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>>> = (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)] diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index c814aa114..7fd230caf 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -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::()) + .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();