fix(ecstore): bound copy-source shard read-ahead (#6663)

This commit is contained in:
cxymds
2026-08-26 21:24:37 +08:00
committed by GitHub
parent a96dd7d289
commit 7c2361757e
13 changed files with 2262 additions and 145 deletions
+571 -9
View File
@@ -35,7 +35,11 @@ use crate::disk::{
health_state::{RuntimeDriveHealthState, get_drive_returning_probe_interval, record_drive_runtime_state},
validate_batch_read_version_item_count,
};
use crate::disk::{disk_store::DiskHealthTracker, error::DiskError, local::ScanGuard};
use crate::disk::{
disk_store::DiskHealthTracker,
error::{DiskError, is_terminal_read_error, terminal_read_error_to_io},
local::ScanGuard,
};
use crate::set_disk::DEFAULT_READ_BUFFER_SIZE;
use bytes::Bytes;
use futures::lock::Mutex;
@@ -324,6 +328,39 @@ where
}
}
/// Mark a terminal fresh-shard recovery failure for adaptive retirement while
/// retaining its typed `DiskError` and original I/O kind. The decoder checks
/// the marker independently of the kind because not-found and transport
/// failures are terminal too, but must not be reported as timeouts.
fn remote_read_error_to_io(error: DiskError) -> io::Error {
terminal_read_error_to_io(error)
}
/// Retire a remote shard after its stream can no longer be trusted. A body
/// error that arrives after the one permitted resume is terminal: retaining
/// the reader would let the next stripe poll an already misaligned stream.
fn remote_terminal_io_error(error: io::Error) -> io::Error {
if is_terminal_read_error(&error) {
return error;
}
terminal_read_error_to_io(DiskError::from(error))
}
fn remote_terminal_message_to_io(message: &'static str) -> io::Error {
terminal_read_error_to_io(DiskError::Io(io::Error::other(message)))
}
fn remote_terminal_eof_to_io() -> io::Error {
terminal_read_error_to_io(DiskError::Io(io::Error::new(
io::ErrorKind::UnexpectedEof,
"remote read ended before requested length",
)))
}
fn remote_terminal_task_error_to_io(error: JoinError) -> io::Error {
terminal_read_error_to_io(DiskError::other(error))
}
struct AbortOnDropTask<T>(JoinHandle<T>);
impl<T> AbortOnDropTask<T> {
@@ -418,6 +455,9 @@ impl RetryingRemoteReader {
impl AsyncRead for RetryingRemoteReader {
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<io::Result<()>> {
if buf.remaining() == 0 {
return Poll::Ready(Ok(()));
}
loop {
// After the absolute cutoff, let initial progress win over a stale fresh-open.
let resume_pending = if let Some(resume) = self.resume.as_mut() {
@@ -431,14 +471,14 @@ impl AsyncRead for RetryingRemoteReader {
Poll::Ready(Ok(Err(error))) => {
self.resume = None;
if self.reader.is_none() {
return Poll::Ready(Err(io::Error::other(error)));
return Poll::Ready(Err(remote_read_error_to_io(error)));
}
continue;
}
Poll::Ready(Err(error)) => {
self.resume = None;
if self.reader.is_none() {
return Poll::Ready(Err(io::Error::other(error)));
return Poll::Ready(Err(remote_terminal_task_error_to_io(error)));
}
continue;
}
@@ -479,6 +519,9 @@ impl AsyncRead for RetryingRemoteReader {
} else {
self.resume = None;
}
} else if produced == 0 && self.request.length != 0 && self.emitted < self.request.length {
self.reader = None;
return Poll::Ready(Err(remote_terminal_eof_to_io()));
}
return Poll::Ready(Ok(()));
}
@@ -494,7 +537,10 @@ impl AsyncRead for RetryingRemoteReader {
self.reader = None;
continue;
}
Poll::Ready(Err(error)) => return Poll::Ready(Err(error)),
Poll::Ready(Err(error)) => {
self.reader = None;
return Poll::Ready(Err(remote_terminal_io_error(error)));
}
}
}
}
@@ -585,21 +631,23 @@ impl rustfs_rio::ChunkReader for RetryingRemoteChunkReader {
Poll::Ready(Ok(Ok(None))) => {
self.resume = None;
if self.reader.is_none() {
return Poll::Ready(Err(io::Error::other("remote resume transport did not provide a chunk reader")));
return Poll::Ready(Err(remote_terminal_message_to_io(
"remote resume transport did not provide a chunk reader",
)));
}
continue;
}
Poll::Ready(Ok(Err(error))) => {
self.resume = None;
if self.reader.is_none() {
return Poll::Ready(Err(io::Error::other(error)));
return Poll::Ready(Err(remote_read_error_to_io(error)));
}
continue;
}
Poll::Ready(Err(error)) => {
self.resume = None;
if self.reader.is_none() {
return Poll::Ready(Err(io::Error::other(error)));
return Poll::Ready(Err(remote_terminal_task_error_to_io(error)));
}
continue;
}
@@ -635,10 +683,26 @@ impl rustfs_rio::ChunkReader for RetryingRemoteChunkReader {
return Poll::Ready(Ok(Some(chunk)));
}
Poll::Ready(Ok(None)) if resume_pending => {
// A clean EOF from the original stream wins when the
// request is unbounded (or has already emitted its full
// bounded length). Waiting for a speculative fresh open
// in that case can turn a successful read into a recovery
// timeout, especially on the read_file/unbounded path.
if self.request.length == 0 || self.emitted >= self.request.length {
self.reader = None;
self.resume = None;
return Poll::Ready(Ok(None));
}
self.reader = None;
continue;
}
Poll::Ready(Ok(None)) => return Poll::Ready(Ok(None)),
Poll::Ready(Ok(None)) => {
if self.request.length != 0 && self.emitted < self.request.length {
self.reader = None;
return Poll::Ready(Err(remote_terminal_eof_to_io()));
}
return Poll::Ready(Ok(None));
}
Poll::Ready(Err(error)) if !self.retried && is_retryable_remote_body_error(&error) => {
self.retried = true;
self.reader = None;
@@ -651,7 +715,10 @@ impl rustfs_rio::ChunkReader for RetryingRemoteChunkReader {
self.reader = None;
continue;
}
Poll::Ready(Err(error)) => return Poll::Ready(Err(error)),
Poll::Ready(Err(error)) => {
self.reader = None;
return Poll::Ready(Err(remote_terminal_io_error(error)));
}
}
}
}
@@ -5154,6 +5221,7 @@ mod tests {
enum ResumeReadStep {
PartialThenReset(Vec<u8>),
Data(Vec<u8>),
Eof,
}
#[derive(Debug, Default)]
@@ -5412,6 +5480,7 @@ mod tests {
struct PendingFreshOpenTransport {
fresh_read_drops: Arc<AtomicUsize>,
fresh_chunk_drops: Arc<AtomicUsize>,
initial_chunk_eof: bool,
}
#[async_trait::async_trait]
@@ -5429,6 +5498,9 @@ mod tests {
}
async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result<Option<rustfs_rio::ChunkReaderBox>> {
if self.initial_chunk_eof {
return Ok(Some(resume_step_chunk_reader(ResumeReadStep::Eof)));
}
Ok(Some(Box::new(ChunkPartialThenErrorReader {
data: None,
error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")),
@@ -5457,6 +5529,70 @@ mod tests {
}
}
#[derive(Debug)]
struct TerminalFreshOpenTransport {
fresh_read_opens: Arc<AtomicUsize>,
fresh_chunk_opens: Arc<AtomicUsize>,
chunk_returns_none: bool,
}
impl TerminalFreshOpenTransport {
fn new(chunk_returns_none: bool) -> Self {
Self {
fresh_read_opens: Arc::new(AtomicUsize::new(0)),
fresh_chunk_opens: Arc::new(AtomicUsize::new(0)),
chunk_returns_none,
}
}
}
#[async_trait::async_trait]
impl InternodeDataTransport for TerminalFreshOpenTransport {
async fn open_read(&self, _request: ReadStreamRequest) -> Result<FileReader> {
Ok(Box::new(PartialThenErrorReader {
cursor: Cursor::new(Vec::new()),
error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")),
}))
}
async fn open_read_fresh(&self, _request: ReadStreamRequest) -> Result<FileReader> {
self.fresh_read_opens.fetch_add(1, Ordering::Relaxed);
Err(DiskError::FileNotFound)
}
async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result<Option<rustfs_rio::ChunkReaderBox>> {
Ok(Some(Box::new(ChunkPartialThenErrorReader {
data: Some(Bytes::from_static(b"x")),
error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")),
})))
}
async fn open_read_chunks_fresh(&self, _request: ReadStreamRequest) -> Result<Option<rustfs_rio::ChunkReaderBox>> {
self.fresh_chunk_opens.fetch_add(1, Ordering::Relaxed);
if self.chunk_returns_none {
Ok(None)
} else {
Err(DiskError::FileNotFound)
}
}
async fn open_write(&self, _request: WriteStreamRequest) -> Result<FileWriter> {
panic!("open_write should not be used in terminal fresh-open tests");
}
async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result<FileReader> {
panic!("open_walk_dir should not be used in terminal fresh-open tests");
}
fn name(&self) -> &'static str {
"terminal-fresh-open-test"
}
fn capabilities(&self) -> InternodeDataTransportCapabilities {
InternodeDataTransportCapabilities::tcp_http()
}
}
fn resume_step_reader(step: ResumeReadStep) -> FileReader {
match step {
ResumeReadStep::PartialThenReset(data) => Box::new(PartialThenErrorReader {
@@ -5464,6 +5600,7 @@ mod tests {
error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")),
}),
ResumeReadStep::Data(data) => Box::new(Cursor::new(data)),
ResumeReadStep::Eof => Box::new(Cursor::new(Vec::new())),
}
}
@@ -5477,6 +5614,7 @@ mod tests {
data: Some(Bytes::from(data)),
error: None,
}),
ResumeReadStep::Eof => Box::new(ChunkPartialThenErrorReader { data: None, error: None }),
}
}
@@ -5553,6 +5691,430 @@ mod tests {
}
}
fn partial_hashed_shard(shard_size: usize) -> (rustfs_utils::HashAlgorithm, Vec<u8>, usize) {
let checksum = rustfs_utils::HashAlgorithm::HighwayHash256S;
let data = vec![0x5a; shard_size];
let hash_bytes = {
let hash = checksum.hash_encode(&data);
hash.as_ref().to_vec()
};
let hash_len = hash_bytes.len();
let encoded_length = hash_len + data.len();
let mut prefix = Vec::with_capacity(hash_len + shard_size / 2);
prefix.extend_from_slice(&hash_bytes);
prefix.extend_from_slice(&data[..shard_size / 2]);
(checksum, prefix, encoded_length)
}
#[test]
fn remote_read_error_conversion_preserves_recovery_classification() {
for disk_error in [DiskError::Timeout, DiskError::SourceStalled] {
let error = remote_read_error_to_io(disk_error);
assert_eq!(error.kind(), std_io::ErrorKind::TimedOut);
}
let error = remote_read_error_to_io(DiskError::Timeout);
assert!(
error
.get_ref()
.and_then(|source| source.downcast_ref::<crate::disk::error::TerminalReadError>())
.is_some()
);
assert!(matches!(DiskError::from(error), DiskError::Timeout));
let error = remote_read_error_to_io(DiskError::SourceStalled);
assert!(matches!(DiskError::from(error), DiskError::SourceStalled));
let error =
remote_read_error_to_io(DiskError::Io(io::Error::new(std_io::ErrorKind::ConnectionReset, "connection reset")));
assert_eq!(error.kind(), std_io::ErrorKind::ConnectionReset);
assert!(crate::disk::error::is_terminal_read_error(&error));
assert!(matches!(DiskError::from(error), DiskError::Io(inner) if inner.kind() == std_io::ErrorKind::ConnectionReset));
}
#[test]
fn remote_reader_zero_capacity_poll_is_a_noop() {
let transport: Arc<dyn InternodeDataTransport> = Arc::new(PendingFreshOpenTransport::default());
let mut reader = RetryingRemoteReader::new_with_timeouts(
Box::new(Cursor::new(b"x".to_vec())),
transport,
resume_request(1),
None,
None,
);
let mut empty = [];
let mut read_buf = ReadBuf::new(&mut empty);
let mut cx = Context::from_waker(std::task::Waker::noop());
assert!(matches!(Pin::new(&mut reader).poll_read(&mut cx, &mut read_buf), Poll::Ready(Ok(()))));
assert!(reader.reader.is_some(), "zero-capacity polls must not retire the remote reader");
let mut output = Vec::new();
futures::executor::block_on(reader.read_to_end(&mut output)).expect("the reader should remain usable");
assert_eq!(output, b"x");
}
#[tokio::test(start_paused = true)]
async fn remote_reader_fresh_open_timeout_preserves_timed_out_kind() {
let transport = Arc::new(PendingFreshOpenTransport::default());
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let mut reader = RetryingRemoteReader::new_with_timeouts(
resume_step_reader(ResumeReadStep::PartialThenReset(Vec::new())),
transport_for_reader,
resume_request(1),
None,
Some(Duration::from_secs(1)),
);
let error = reader
.read_to_end(&mut Vec::new())
.await
.expect_err("a hung fresh open must surface its recovery timeout");
assert_eq!(error.kind(), std_io::ErrorKind::TimedOut);
assert!(
error
.get_ref()
.and_then(|source| source.downcast_ref::<crate::disk::error::TerminalReadError>())
.is_some()
);
assert!(matches!(DiskError::from(error), DiskError::Timeout));
assert_eq!(transport.fresh_read_drops.load(Ordering::Relaxed), 1);
}
#[tokio::test(start_paused = true)]
async fn remote_chunk_reader_fresh_open_timeout_preserves_timed_out_kind() {
let transport = Arc::new(PendingFreshOpenTransport::default());
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let mut reader = RetryingRemoteChunkReader::new_with_timeouts(
resume_step_chunk_reader(ResumeReadStep::PartialThenReset(b"x".to_vec())),
transport_for_reader,
resume_request(2),
None,
Some(Duration::from_secs(1)),
);
let error = reader
.read_to_end(&mut Vec::new())
.await
.expect_err("a hung fresh chunk open must surface its recovery timeout");
assert_eq!(error.kind(), std_io::ErrorKind::TimedOut);
assert!(
error
.get_ref()
.and_then(|source| source.downcast_ref::<crate::disk::error::TerminalReadError>())
.is_some()
);
assert!(matches!(DiskError::from(error), DiskError::Timeout));
assert_eq!(transport.fresh_chunk_drops.load(Ordering::Relaxed), 1);
}
#[tokio::test(start_paused = true)]
async fn remote_chunk_reader_unbounded_clean_eof_wins_over_speculative_resume() {
let transport = Arc::new(PendingFreshOpenTransport {
initial_chunk_eof: true,
..PendingFreshOpenTransport::default()
});
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let mut reader = RetryingRemoteChunkReader::new_with_timeouts(
resume_step_chunk_reader(ResumeReadStep::Eof),
transport_for_reader,
resume_request(0),
Some(Duration::ZERO),
Some(Duration::from_secs(1)),
);
let mut output = Vec::new();
reader
.read_to_end(&mut output)
.await
.expect("clean EOF from an unbounded original stream should finish the read");
assert!(output.is_empty());
// The executor may abort the speculative task before it is first
// polled, in which case the pending-open future never constructs its
// drop probe. The dedicated drop-cancellation tests cover the
// already-polled case; this regression only needs to establish that a
// clean unbounded EOF is not converted into a recovery timeout.
}
#[tokio::test(start_paused = true)]
async fn remote_reader_fresh_open_non_timeout_error_is_retired_from_adaptive_decode() {
let transport = Arc::new(TerminalFreshOpenTransport::new(false));
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteReader::new_with_timeouts(
resume_step_reader(ResumeReadStep::PartialThenReset(Vec::new())),
transport_for_reader,
resume_request(8),
None,
Some(Duration::from_secs(1)),
);
let shard = BitrotReader::new(ShardReader::Stream(Box::new(retry)), 8, rustfs_utils::HashAlgorithm::None, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, 8), 0, 16);
let (_, first_errors) = parallel.read().await;
assert!(matches!(first_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(transport.fresh_read_opens.load(Ordering::Relaxed), 1);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(transport.fresh_read_opens.load(Ordering::Relaxed), 1);
}
#[tokio::test(start_paused = true)]
async fn remote_chunk_reader_missing_fresh_reader_is_retired_from_adaptive_decode() {
let transport = Arc::new(TerminalFreshOpenTransport::new(true));
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteChunkReader::new_with_timeouts(
resume_step_chunk_reader(ResumeReadStep::PartialThenReset(b"x".to_vec())),
transport_for_reader,
resume_request(8),
None,
Some(Duration::from_secs(1)),
);
let shard = BitrotReader::new(ShardReader::Chunked(Box::new(retry)), 8, rustfs_utils::HashAlgorithm::None, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, 8), 0, 16);
let (_, first_errors) = parallel.read().await;
assert!(
matches!(first_errors.first().and_then(Option::as_ref), Some(DiskError::Io(error)) if error.kind() == std_io::ErrorKind::Other),
"unexpected first errors: {first_errors:?}"
);
assert_eq!(transport.fresh_chunk_opens.load(Ordering::Relaxed), 1);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(transport.fresh_chunk_opens.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn remote_reader_fresh_open_short_eof_is_retired_from_adaptive_decode() {
// A successful fresh open can still end before the bounded request.
// Treat that as terminal immediately so the next stripe does not poll
// an already exhausted reader and defer the failure to BitrotReader.
let transport = Arc::new(ResumeTransport::with_read_steps(vec![ResumeReadStep::Eof]));
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteReader::new_with_timeouts(
resume_step_reader(ResumeReadStep::PartialThenReset(b"01".to_vec())),
transport_for_reader,
resume_request(7),
None,
None,
);
let shard = BitrotReader::new(ShardReader::Stream(Box::new(retry)), 7, rustfs_utils::HashAlgorithm::None, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, 7), 0, 14);
let (_, first_errors) = parallel.read().await;
assert!(matches!(
first_errors.first().and_then(Option::as_ref),
Some(DiskError::Io(error)) if error.kind() == std_io::ErrorKind::UnexpectedEof
));
assert_eq!(
transport
.fresh_read_requests
.lock()
.expect("fresh read request lock should not be poisoned")
.len(),
1
);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(
transport
.fresh_read_requests
.lock()
.expect("fresh read request lock should not be poisoned")
.len(),
1
);
}
#[tokio::test]
async fn remote_chunk_reader_fresh_open_short_eof_is_retired_from_adaptive_decode() {
let transport = Arc::new(ResumeTransport::with_chunk_steps(vec![ResumeReadStep::Eof]));
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteChunkReader::new_with_timeouts(
resume_step_chunk_reader(ResumeReadStep::PartialThenReset(b"01".to_vec())),
transport_for_reader,
resume_request(7),
None,
None,
);
let shard = BitrotReader::new(ShardReader::Chunked(Box::new(retry)), 7, rustfs_utils::HashAlgorithm::None, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, 7), 0, 14);
let (_, first_errors) = parallel.read().await;
assert!(matches!(
first_errors.first().and_then(Option::as_ref),
Some(DiskError::Io(error)) if error.kind() == std_io::ErrorKind::UnexpectedEof
));
assert_eq!(
transport
.fresh_chunk_requests
.lock()
.expect("fresh chunk request lock should not be poisoned")
.len(),
1
);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(
transport
.fresh_chunk_requests
.lock()
.expect("fresh chunk request lock should not be poisoned")
.len(),
1
);
}
#[tokio::test]
async fn remote_reader_second_body_reset_is_retired_from_adaptive_decode() {
// The first connection emits a prefix, the one permitted fresh
// connection emits another prefix, and then resets again. The second
// reset must retire the reader so the next stripe cannot consume a
// misaligned stream.
let transport = Arc::new(ResumeTransport::with_read_steps(vec![ResumeReadStep::PartialThenReset(b"23".to_vec())]));
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteReader::new_with_timeouts(
resume_step_reader(ResumeReadStep::PartialThenReset(b"01".to_vec())),
transport_for_reader,
resume_request(7),
None,
None,
);
let shard = BitrotReader::new(ShardReader::Stream(Box::new(retry)), 7, rustfs_utils::HashAlgorithm::None, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, 7), 0, 14);
let (_, first_errors) = parallel.read().await;
assert!(matches!(
first_errors.first().and_then(Option::as_ref),
Some(DiskError::Io(error)) if error.kind() == std_io::ErrorKind::ConnectionReset
));
assert_eq!(
transport
.fresh_read_requests
.lock()
.expect("fresh read request lock should not be poisoned")
.len(),
1
);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(
transport
.fresh_read_requests
.lock()
.expect("fresh read request lock should not be poisoned")
.len(),
1
);
}
#[tokio::test]
async fn remote_chunk_reader_second_body_reset_is_retired_from_adaptive_decode() {
let transport = Arc::new(ResumeTransport::with_chunk_steps(vec![ResumeReadStep::PartialThenReset(b"23".to_vec())]));
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteChunkReader::new_with_timeouts(
resume_step_chunk_reader(ResumeReadStep::PartialThenReset(b"01".to_vec())),
transport_for_reader,
resume_request(7),
None,
None,
);
let shard = BitrotReader::new(ShardReader::Chunked(Box::new(retry)), 7, rustfs_utils::HashAlgorithm::None, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, 7), 0, 14);
let (_, first_errors) = parallel.read().await;
assert!(matches!(
first_errors.first().and_then(Option::as_ref),
Some(DiskError::Io(error)) if error.kind() == std_io::ErrorKind::ConnectionReset
));
assert_eq!(
transport
.fresh_chunk_requests
.lock()
.expect("fresh chunk request lock should not be poisoned")
.len(),
1
);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
assert_eq!(
transport
.fresh_chunk_requests
.lock()
.expect("fresh chunk request lock should not be poisoned")
.len(),
1
);
}
#[tokio::test(start_paused = true)]
#[serial]
async fn remote_reader_hashed_timeout_is_retired_from_adaptive_decode() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_DISK_READ_TIMEOUT, Some("5"))], async {
const SHARD_SIZE: usize = 64;
let (checksum, encoded_prefix, encoded_length) = partial_hashed_shard(SHARD_SIZE);
let transport = Arc::new(PendingFreshOpenTransport::default());
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteReader::new_with_timeouts(
resume_step_reader(ResumeReadStep::PartialThenReset(encoded_prefix)),
transport_for_reader,
resume_request(encoded_length),
None,
Some(Duration::from_secs(1)),
);
let shard = BitrotReader::new(ShardReader::Stream(Box::new(retry)), SHARD_SIZE, checksum, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, SHARD_SIZE), 0, SHARD_SIZE * 2);
let (first_buffers, first_errors) = parallel.read().await;
assert!(first_buffers.first().and_then(Option::as_ref).is_none());
assert!(matches!(first_errors.first().and_then(Option::as_ref), Some(DiskError::Timeout)));
assert_eq!(transport.fresh_read_drops.load(Ordering::Relaxed), 1);
// A TimedOut error retires the dead slot. The next stripe therefore
// reports the slot as unavailable instead of polling a reader whose
// fresh connection already timed out.
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
})
.await;
}
#[tokio::test(start_paused = true)]
#[serial]
async fn remote_chunk_reader_hashed_timeout_is_retired_from_adaptive_decode() {
temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_DISK_READ_TIMEOUT, Some("5"))], async {
const SHARD_SIZE: usize = 64;
let (checksum, encoded_prefix, encoded_length) = partial_hashed_shard(SHARD_SIZE);
let transport = Arc::new(PendingFreshOpenTransport::default());
let transport_for_reader: Arc<dyn InternodeDataTransport> = transport.clone();
let retry = RetryingRemoteChunkReader::new_with_timeouts(
resume_step_chunk_reader(ResumeReadStep::PartialThenReset(encoded_prefix)),
transport_for_reader,
resume_request(encoded_length),
None,
Some(Duration::from_secs(1)),
);
let shard = BitrotReader::new(ShardReader::Chunked(Box::new(retry)), SHARD_SIZE, checksum, false);
let mut parallel = ParallelReader::new(vec![Some(shard)], Erasure::new(1, 0, SHARD_SIZE), 0, SHARD_SIZE * 2);
let (first_buffers, first_errors) = parallel.read().await;
assert!(first_buffers.first().and_then(Option::as_ref).is_none());
assert!(matches!(first_errors.first().and_then(Option::as_ref), Some(DiskError::Timeout)));
assert_eq!(transport.fresh_chunk_drops.load(Ordering::Relaxed), 1);
let (_, second_errors) = parallel.read().await;
assert!(matches!(second_errors.first().and_then(Option::as_ref), Some(DiskError::FileNotFound)));
})
.await;
}
#[tokio::test]
async fn remote_reader_resumes_from_emitted_bytes_without_duplicates() {
let transport = Arc::new(ResumeTransport::with_read_steps(vec![ResumeReadStep::Data(b"456789".to_vec())]));
+112
View File
@@ -23,6 +23,16 @@ pub type Result<T> = core::result::Result<T, Error>;
const METACACHE_OUTPUT_STREAM_CLOSED: &str = "metacache output stream closed";
/// Marker carried by a shard-read `io::Error` when the underlying reader can
/// no longer be realigned after a fresh remote open failed. The marker is
/// deliberately separate from the `ErrorKind`: a terminal read must retire
/// its reader, while its original typed disk error and I/O kind still need to
/// survive quorum/error mapping.
#[derive(Debug)]
pub(crate) struct TerminalReadError {
source: DiskError,
}
// DiskError == StorageErr
#[derive(Debug, thiserror::Error)]
pub enum DiskError {
@@ -168,6 +178,67 @@ pub enum DiskError {
RemoteClientUnavailable(String),
}
impl TerminalReadError {
pub(crate) fn new(source: DiskError) -> Self {
Self { source }
}
fn into_source(self) -> DiskError {
self.source
}
}
impl std::fmt::Display for TerminalReadError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.source.fmt(f)
}
}
impl StdError for TerminalReadError {
fn source(&self) -> Option<&(dyn StdError + 'static)> {
Some(&self.source)
}
}
fn classify_internode_missing_error(error: &InternodeHttpError) -> Option<DiskError> {
if error.is_remote_file_not_found() {
return Some(DiskError::FileNotFound);
}
if error.is_remote_volume_not_found() {
return Some(DiskError::VolumeNotFound);
}
None
}
/// Wrap a terminal shard-read failure without changing its typed
/// classification. Timeout-like disk errors retain `TimedOut`; other errors
/// retain their inner I/O kind or use `Other` when no more specific kind exists.
pub(crate) fn terminal_read_error_to_io(error: DiskError) -> io::Error {
let kind = match &error {
DiskError::Io(inner) => inner.kind(),
DiskError::SourceStalled | DiskError::Timeout => io::ErrorKind::TimedOut,
DiskError::DiskNotFound
| DiskError::FileNotFound
| DiskError::FileVersionNotFound
| DiskError::PathNotFound
| DiskError::VolumeNotFound => io::ErrorKind::NotFound,
DiskError::DiskAccessDenied | DiskError::FileAccessDenied | DiskError::VolumeAccessDenied => {
io::ErrorKind::PermissionDenied
}
DiskError::DiskFull => io::ErrorKind::StorageFull,
DiskError::FileCorrupt | DiskError::PartMissingOrCorrupt | DiskError::BitrotHashAlgoInvalid => io::ErrorKind::InvalidData,
_ => io::ErrorKind::Other,
};
io::Error::new(kind, TerminalReadError::new(error))
}
/// Whether an I/O error marks a shard reader as terminal for adaptive decode.
pub(crate) fn is_terminal_read_error(error: &io::Error) -> bool {
error
.get_ref()
.is_some_and(|source| source.downcast_ref::<TerminalReadError>().is_some())
}
impl From<crate::erasure::coding::ErasureConstructionError> for DiskError {
fn from(error: crate::erasure::coding::ErasureConstructionError) -> Self {
Self::Io(error.into_io_error())
@@ -344,6 +415,21 @@ impl From<std::io::Error> for DiskError {
return DiskError::VolumeNotFound;
}
}
let e = match e.downcast::<TerminalReadError>() {
Ok(terminal_error) => {
let source = terminal_error.into_source();
if let DiskError::Io(io_error) = &source
&& let Some(internode_error) = io_error
.get_ref()
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
&& let Some(classified) = classify_internode_missing_error(internode_error)
{
return classified;
}
return source;
}
Err(e) => e,
};
match e.downcast::<DiskError>() {
Ok(disk_error) => disk_error,
// Mirror `From<io::Error> for StorageError`: a StorageError boxed
@@ -679,6 +765,32 @@ mod tests {
use super::*;
use std::collections::HashMap;
#[test]
fn terminal_read_error_preserves_kind_and_disk_classification() {
let timeout = terminal_read_error_to_io(DiskError::Timeout);
assert_eq!(timeout.kind(), io::ErrorKind::TimedOut);
assert!(is_terminal_read_error(&timeout));
assert!(matches!(DiskError::from(timeout), DiskError::Timeout));
let missing = terminal_read_error_to_io(DiskError::FileNotFound);
assert_eq!(missing.kind(), io::ErrorKind::NotFound);
assert!(is_terminal_read_error(&missing));
assert!(matches!(DiskError::from(missing), DiskError::FileNotFound));
let reset = terminal_read_error_to_io(DiskError::Io(io::Error::new(io::ErrorKind::ConnectionReset, "connection reset")));
assert_eq!(reset.kind(), io::ErrorKind::ConnectionReset);
assert!(is_terminal_read_error(&reset));
assert!(matches!(DiskError::from(reset), DiskError::Io(error) if error.kind() == io::ErrorKind::ConnectionReset));
for (remote_error, expected) in [
(rustfs_rio::new_test_remote_file_not_found_http_io_error(), DiskError::FileNotFound),
(rustfs_rio::new_test_remote_volume_not_found_http_io_error(), DiskError::VolumeNotFound),
] {
let wrapped = terminal_read_error_to_io(DiskError::Io(remote_error));
assert_eq!(DiskError::from(wrapped), expected);
}
}
#[test]
fn other_preserves_erasure_construction_source_chain() {
use crate::erasure::coding::ErasureConstructionError;
File diff suppressed because it is too large Load Diff
+5 -10
View File
@@ -346,16 +346,11 @@ impl AsyncRead for DeferredObjectReader {
}
fn disk_error_to_io_error(err: DiskError) -> io::Error {
let kind = match err {
DiskError::Timeout | DiskError::SourceStalled => io::ErrorKind::TimedOut,
DiskError::DiskNotFound | DiskError::FileNotFound | DiskError::FileVersionNotFound | DiskError::PathNotFound => {
io::ErrorKind::NotFound
}
DiskError::FileCorrupt | DiskError::PartMissingOrCorrupt | DiskError::BitrotHashAlgoInvalid => io::ErrorKind::InvalidData,
DiskError::Io(io_err) => return io_err,
_ => io::ErrorKind::Other,
};
io::Error::new(kind, err.to_string())
// Keep the typed disk error attached to deferred-reader failures. The
// decoder uses the marker to retire a stream that can no longer be
// realigned, while quorum reduction still sees Timeout/NotFound instead
// of an opaque `DiskError::Io` wrapper.
crate::disk::error::terminal_read_error_to_io(err)
}
async fn open_disk_reader(
@@ -1517,6 +1517,7 @@ pub(in crate::set_disk) async fn submit_read_repair_heal_with_submitter(
}
pub(in crate::set_disk) type ObjectBitrotReader = BitrotReader<ShardReader>;
pub(in crate::set_disk) type DeferredReaderReopener = crate::erasure::coding::decode::DeferredReaderReopener<ShardReader>;
pub(in crate::set_disk) type BitrotReaderTask<'a> =
Pin<Box<dyn Future<Output = (usize, std::result::Result<Option<ObjectBitrotReader>, DiskError>)> + Send + 'a>>;
@@ -1533,6 +1534,10 @@ pub(in crate::set_disk) struct BitrotReaderSetup {
/// 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>>,
/// Factories for a fresh, stripe-aligned parity reader. CopySource hedges
/// use these disposable readers so an abandoned hedge leaves the original
/// deferred reserve untouched.
pub(in crate::set_disk) deferred_reopeners: Vec<Option<DeferredReaderReopener>>,
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>,
@@ -1595,6 +1600,16 @@ pub(in crate::set_disk) fn get_bitrot_reader_setup_strategy(
mode: BitrotReaderSetupMode,
prefer_data_blocks_first: bool,
) -> BitrotReaderSetupStrategy {
// CopyObject holds the source reader behind a backpressured destination.
// Keep its setup demand-bound even when an operator has retained the
// legacy all-shards environment setting for ordinary GETs.
if matches!(
crate::set_disk::get_object_read_policy(),
crate::set_disk::GetObjectReadPolicy::CopySource
) {
return BitrotReaderSetupStrategy::DataBlocksFirst;
}
match mode {
BitrotReaderSetupMode::ReadQuorum
if prefer_data_blocks_first
@@ -1620,6 +1635,7 @@ impl BitrotReaderSetup {
Self {
readers: (0..shards).map(|_| None).collect(),
deferred_stripe_handles: (0..shards).map(|_| None).collect(),
deferred_reopeners: (0..shards).map(|_| None).collect(),
errors: vec![Some(DiskError::DiskNotFound); shards],
scheduled: vec![false; shards],
attempted: vec![false; shards],
@@ -1814,6 +1830,41 @@ pub(in crate::set_disk) fn next_unscheduled_reader_index(
.find(|idx| !setup.scheduled[*idx])
}
/// Build a cloneable opener for an unopened deferred shard. The returned
/// reader is aligned to the requested stripe before its first poll, while the
/// source reader created during setup remains untouched as a reserve.
#[allow(clippy::too_many_arguments)]
fn deferred_reader_reopener(
inline_data: Option<Bytes>,
disk: Option<DiskStore>,
bucket: &str,
path: &str,
read_offset: usize,
read_length: usize,
shard_size: usize,
checksum_algo: HashAlgorithm,
skip_verify_bitrot: bool,
use_mmap_read: bool,
) -> DeferredReaderReopener {
let bucket = bucket.to_owned();
let path = path.to_owned();
Arc::new(move |stripe_index| {
let (reader, handle) = create_deferred_bitrot_reader_with_stripe_handle(
inline_data.clone(),
disk.clone(),
&bucket,
&path,
read_offset,
read_length,
shard_size,
checksum_algo.clone(),
skip_verify_bitrot,
use_mmap_read,
);
handle.advance_stripes(stripe_index).then_some(reader)
})
}
#[allow(clippy::too_many_arguments)]
pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
setup: &mut BitrotReaderSetup,
@@ -1836,6 +1887,15 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
return;
}
// Only CopySource uses disposable, stripe-aligned reopeners. Ordinary GET
// readers use the existing deferred handle and should not retain one
// heap-allocated closure (plus cloned path/disk state) for every parity
// slot.
let copy_source_demand_bound = matches!(
crate::set_disk::get_object_read_policy(),
crate::set_disk::GetObjectReadPolicy::CopySource
);
for idx in 0..disks.len() {
if setup.attempted[idx] {
continue;
@@ -1849,6 +1909,20 @@ 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}");
let reopener = copy_source_demand_bound.then(|| {
deferred_reader_reopener(
inline_data.clone(),
disk.clone(),
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,
@@ -1862,6 +1936,7 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
use_mmap_read,
);
setup.retain_deferred_reader(idx, reader, stripe_handle);
setup.deferred_reopeners[idx] = reopener;
}
// With the data-shards-only lockstep gate on (backlog#923), the GET decode
@@ -1887,6 +1962,20 @@ 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}");
let reopener = copy_source_demand_bound.then(|| {
deferred_reader_reopener(
inline_data.clone(),
disk.clone(),
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,
@@ -1901,6 +1990,7 @@ pub(in crate::set_disk) fn fill_deferred_bitrot_readers(
);
setup.readers[idx] = Some(reader);
setup.deferred_stripe_handles[idx] = Some(stripe_handle);
setup.deferred_reopeners[idx] = reopener;
}
}
@@ -2210,6 +2300,10 @@ pub(in crate::set_disk) async fn create_bitrot_readers_until_quorum_with_prefere
let strategy = get_bitrot_reader_setup_strategy(mode, prefer_data_blocks_first);
if use_mmap_read
&& !matches!(
crate::set_disk::get_object_read_policy(),
crate::set_disk::GetObjectReadPolicy::CopySource
)
&& let Some(mut setup) = try_create_bitrot_readers_via_batch_pread(
files,
disks,
+95
View File
@@ -792,6 +792,65 @@ const ENV_RUSTFS_GET_METADATA_SLOWTAIL_FAULT_OBJECT_PREFIX: &str = "RUSTFS_GET_M
const ENV_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH: &str = "RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH";
const DEFAULT_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH: bool = true;
/// Identifies the caller's read contract for policies that are deliberately
/// narrower than the storage API's ordinary GET contract.
///
/// Server-side copy consumes a source reader while a destination writer is
/// applying backpressure. Its source read must not speculatively open the
/// next multipart part: those extra shard streams can share an internode H2
/// connection with the current part and starve the lockstep decoder. Keep
/// this context internal so the public `ObjectOptions` and storage traits do
/// not acquire a copy-only field.
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub(crate) enum GetObjectReadPolicy {
#[default]
Default,
CopySource,
}
impl GetObjectReadPolicy {
pub(crate) const fn allows_multipart_setup_prefetch(self) -> bool {
matches!(self, Self::Default)
}
}
tokio::task_local! {
static GET_OBJECT_READ_POLICY: GetObjectReadPolicy;
static GET_OBJECT_READ_CANCELLATION: tokio_util::sync::CancellationToken;
}
pub(crate) fn get_object_read_policy() -> GetObjectReadPolicy {
GET_OBJECT_READ_POLICY.try_with(|policy| *policy).unwrap_or_default()
}
pub(crate) async fn with_get_object_read_policy<F>(policy: GetObjectReadPolicy, future: F) -> F::Output
where
F: std::future::Future,
{
let decode_policy = match policy {
GetObjectReadPolicy::Default => crate::erasure::coding::decode::DecodeReadPolicy::Default,
GetObjectReadPolicy::CopySource => crate::erasure::coding::decode::DecodeReadPolicy::DemandBound,
};
crate::erasure::coding::decode::with_decode_read_policy(decode_policy, GET_OBJECT_READ_POLICY.scope(policy, future)).await
}
/// Return the request-owned cancellation token for a copy source, when one is
/// installed. The token is read before the detached legacy producer is spawned;
/// Tokio task-local values do not cross that spawn boundary on their own.
pub(crate) fn get_object_read_cancellation() -> Option<tokio_util::sync::CancellationToken> {
GET_OBJECT_READ_CANCELLATION.try_with(|token| token.clone()).ok()
}
pub(crate) async fn with_get_object_read_cancellation<F>(
cancellation: tokio_util::sync::CancellationToken,
future: F,
) -> F::Output
where
F: std::future::Future,
{
GET_OBJECT_READ_CANCELLATION.scope(cancellation, future).await
}
static OBJECT_LOCK_DIAG_ENABLED: OnceLock<bool> = OnceLock::new();
mod core;
@@ -2296,6 +2355,7 @@ enum GetCodecStreamingFallbackReason {
InvalidMinSize,
ReadQuorumNotSafe,
MultipartPartLimit,
CopySourceDemandBound,
}
impl GetCodecStreamingFallbackReason {
@@ -2317,6 +2377,7 @@ impl GetCodecStreamingFallbackReason {
Self::InvalidMinSize => "invalid_min_size",
Self::ReadQuorumNotSafe => "read_quorum_not_safe",
Self::MultipartPartLimit => "multipart_part_limit",
Self::CopySourceDemandBound => "copy_source_demand_bound",
}
}
}
@@ -2623,6 +2684,17 @@ fn get_codec_streaming_reader_gate(
prefer_data_blocks_first_reader_setup: false,
};
}
if matches!(get_object_read_policy(), GetObjectReadPolicy::CopySource) {
// The codec reader has its own bounded fill worker. It may still
// request an additional stripe for a plain single-part object even
// when multipart setup prefetch is disabled, so copy sources use the
// legacy demand-bound reader for every object class.
return GetCodecStreamingGate {
object_class,
decision: GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::CopySourceDemandBound),
prefer_data_blocks_first_reader_setup: false,
};
}
if !config.rollout.is_opted_in() {
return GetCodecStreamingGate {
object_class,
@@ -5912,6 +5984,29 @@ mod tests {
use tokio::fs;
use tokio::io::AsyncReadExt;
#[tokio::test]
async fn copy_source_read_policy_is_scoped_and_demand_bound() {
assert_eq!(get_object_read_policy(), GetObjectReadPolicy::Default);
assert!(GetObjectReadPolicy::Default.allows_multipart_setup_prefetch());
assert!(!GetObjectReadPolicy::CopySource.allows_multipart_setup_prefetch());
with_get_object_read_policy(GetObjectReadPolicy::CopySource, async {
assert_eq!(get_object_read_policy(), GetObjectReadPolicy::CopySource);
assert!(!get_object_read_policy().allows_multipart_setup_prefetch());
assert_eq!(
crate::erasure::coding::decode::decode_read_policy(),
crate::erasure::coding::decode::DecodeReadPolicy::DemandBound
);
})
.await;
assert_eq!(get_object_read_policy(), GetObjectReadPolicy::Default);
assert_eq!(
crate::erasure::coding::decode::decode_read_policy(),
crate::erasure::coding::decode::DecodeReadPolicy::Default
);
}
#[test]
fn complete_part_error_maps_confirmed_missing_to_invalid_part() {
for err in ["file not found", "Specified part could not be found", "part.7 not found"] {
+102 -20
View File
@@ -61,6 +61,8 @@ use http::HeaderValue;
use rustfs_utils::path::decode_dir_object;
use std::future::Future;
use std::sync::OnceLock;
use std::task::{Context, Poll};
use tokio::io::{AsyncRead, ReadBuf};
use tokio_util::sync::CancellationToken;
const OLD_DATA_CLEANUP_RECEIPT_FILE: &str = ".rustfs-old-data-cleanup-receipt.json";
@@ -1120,6 +1122,37 @@ where
Ok((reader, offset, length))
}
/// Cancels a detached legacy GET producer when its consumer is dropped.
///
/// The producer owns the shard readers and the object read lock, while the
/// consumer owns only the duplex read half. Closing that half eventually
/// unblocks a writer, but can leave a producer stuck in reader setup or remote
/// recovery until a lower-level timeout fires. This small boundary wrapper
/// provides an explicit cancellation signal without changing the public
/// `GetObjectReader` shape.
struct ProducerCancellationReader<R> {
inner: R,
cancellation: CancellationToken,
}
impl<R> ProducerCancellationReader<R> {
fn new(inner: R, cancellation: CancellationToken) -> Self {
Self { inner, cancellation }
}
}
impl<R: AsyncRead + Unpin> AsyncRead for ProducerCancellationReader<R> {
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
Pin::new(&mut self.inner).poll_read(cx, buf)
}
}
impl<R> Drop for ProducerCancellationReader<R> {
fn drop(&mut self) {
self.cancellation.cancel();
}
}
fn data_read_metadata_early_stop_request_shape_allowed(range: &Option<HTTPRangeSpec>, opts: &ObjectOptions) -> bool {
range.is_none()
&& opts.part_number.is_none()
@@ -1907,12 +1940,25 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
// lookup on the streaming miss path (ODC-16).
reader.body_source = body_source;
// The producer is otherwise detached from the returned reader. Tie its
// lifetime to the source stream so a cancelled copy (or an abandoned
// GET) releases in-flight shard opens, response bodies, and the read
// lock immediately instead of waiting for a disk timeout.
let producer_cancellation = crate::set_disk::get_object_read_cancellation();
if let Some(cancellation) = producer_cancellation.as_ref() {
reader.stream = Box::new(ProducerCancellationReader::new(reader.stream, cancellation.clone()));
}
// let disks = disks.clone();
let bucket = bucket.to_owned();
let object = object.to_owned();
let set_index = self.set_index;
let pool_index = self.pool_index;
let skip_verify = opts.skip_verify_bitrot;
// The producer runs in a separate Tokio task, so carry the caller's
// read policy across the task boundary explicitly. Tokio task-local
// values are not inherited by spawned tasks.
let read_policy = crate::set_disk::get_object_read_policy();
let erasure_cache = Arc::clone(&self.erasure_cache);
let (fi, files, disks) = snapshot.into_owned();
tokio::spawn(async move {
@@ -1922,26 +1968,40 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
// `get_object_with_fileinfo` also waits on `writer`, so an outer timeout
// would incorrectly treat downstream backpressure as disk-read latency.
// Disk read timeouts must be enforced at the actual disk I/O operations.
let producer_result = Self::get_object_with_fileinfo(
&bucket,
&object,
erasure_cache,
offset,
length,
&mut writer,
fi,
files,
&disks,
set_index,
pool_index,
skip_verify,
false,
GET_OBJECT_PATH_LEGACY_DUPLEX,
object_class.as_str(),
size_bucket,
)
.await;
if let Err(e) = &producer_result {
let producer_result = tokio::select! {
biased;
result = crate::set_disk::with_get_object_read_policy(
read_policy,
Self::get_object_with_fileinfo(
&bucket,
&object,
erasure_cache,
offset,
length,
&mut writer,
fi,
files,
&disks,
set_index,
pool_index,
skip_verify,
false,
GET_OBJECT_PATH_LEGACY_DUPLEX,
object_class.as_str(),
size_bucket,
),
) => result,
_ = async {
if let Some(cancellation) = producer_cancellation.as_ref() {
cancellation.cancelled().await;
} else {
std::future::pending::<()>().await;
}
} => Err(Error::OperationCanceled),
};
if let Err(e) = &producer_result
&& !matches!(e, Error::OperationCanceled)
{
let reason = classify_storage_error(e);
if reason == GetObjectFailureReason::DownstreamClosed {
debug!(
@@ -3794,6 +3854,28 @@ mod legacy_duplex_producer_reader_tests {
assert_eq!(out, b"complete");
}
#[tokio::test]
async fn producer_cancellation_reader_cancels_pending_producer_on_drop() {
let cancellation = CancellationToken::new();
let producer_cancellation = cancellation.clone();
let producer = tokio::spawn(async move {
tokio::select! {
_ = producer_cancellation.cancelled() => true,
_ = std::future::pending::<()>() => false,
}
});
let reader = ProducerCancellationReader::new(tokio::io::empty(), cancellation);
drop(reader);
assert!(
tokio::time::timeout(std::time::Duration::from_secs(1), producer)
.await
.expect("dropping the consumer should cancel the producer promptly")
.expect("producer task should not panic")
);
}
#[tokio::test]
async fn legacy_duplex_reader_ignores_zero_capacity_read_buf() {
let (mut writer, reader) = tokio::io::duplex(64);
+76 -4
View File
@@ -792,7 +792,7 @@ impl SetDisks {
let use_mmap_read = object_mmap_read_enabled();
let files = Arc::new(files);
let disks = Arc::new(disks);
let prefetch_enabled = is_multipart_reader_setup_prefetch_enabled();
let prefetch_enabled = multipart_reader_setup_prefetch_enabled(get_object_read_policy());
let mut prefetched: Option<(usize, PrefetchedReaderSetup)> = None;
let mut total_read = 0;
@@ -1065,8 +1065,9 @@ impl SetDisks {
let unattempted_data_shards = !reader_setup.data_shards_attempted(erasure.data_shards);
let readers = reader_setup.readers;
let deferred_stripe_handles = reader_setup.deferred_stripe_handles;
let deferred_reopeners = reader_setup.deferred_reopeners;
let (written, err) = erasure
.decode_with_stripe_handles(
.decode_with_stripe_handles_and_reopeners(
writer,
readers,
part_offset,
@@ -1074,6 +1075,7 @@ impl SetDisks {
part_size,
read_costs,
deferred_stripe_handles,
deferred_reopeners,
)
.await;
let decode_elapsed = decode_stage_start.elapsed();
@@ -1476,6 +1478,7 @@ impl SetDisks {
erasure.clone(),
reader_setup.readers,
reader_setup.deferred_stripe_handles,
reader_setup.deferred_reopeners,
read_costs,
part_offset,
part_length,
@@ -1488,6 +1491,7 @@ impl SetDisks {
let readers = reader_setup.readers;
let deferred_stripe_handles = reader_setup.deferred_stripe_handles;
let deferred_reopeners = reader_setup.deferred_reopeners;
let source = if let Some(read_costs) = read_costs {
coding::decode::ParallelReader::new_with_metrics_path_read_costs_and_reconstruction_verification(
readers,
@@ -1506,7 +1510,8 @@ impl SetDisks {
Some(metrics_path),
)
}
.with_deferred_parity_handles(deferred_stripe_handles);
.with_deferred_parity_handles(deferred_stripe_handles)
.with_deferred_parity_reopeners(deferred_reopeners);
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)?;
@@ -1535,6 +1540,10 @@ fn multipart_part_checksum_algo(fi: &FileInfo, part_number: usize) -> HashAlgori
}
}
fn multipart_reader_setup_prefetch_enabled(policy: GetObjectReadPolicy) -> bool {
policy.allows_multipart_setup_prefetch() && is_multipart_reader_setup_prefetch_enabled()
}
/// Run one part's bitrot reader setup and measure its wall-clock duration.
///
/// Shared by the synchronous path and the prefetch task in
@@ -1762,10 +1771,12 @@ impl Drop for LazyMultipartCodecStreamingReader {
/// background task drives the decode into the write half while the returned
/// reader drains the read half. No extra file descriptors are opened — the
/// readers are moved in from the setup that just ran.
#[allow(clippy::too_many_arguments)]
fn build_legacy_per_part_fallback_reader(
erasure: coding::Erasure,
readers: Vec<Option<ObjectBitrotReader>>,
deferred_stripe_handles: Vec<Option<DeferredReaderStripeHandle>>,
deferred_reopeners: Vec<Option<DeferredReaderReopener>>,
read_costs: Option<Vec<ShardReadCost>>,
part_offset: usize,
part_length: usize,
@@ -1775,7 +1786,7 @@ fn build_legacy_per_part_fallback_reader(
let (read_half, mut write_half) = tokio::io::duplex(buffer);
let decode = tokio::spawn(async move {
let (_written, err) = erasure
.decode_with_stripe_handles(
.decode_with_stripe_handles_and_reopeners(
&mut write_half,
readers,
part_offset,
@@ -1783,6 +1794,7 @@ fn build_legacy_per_part_fallback_reader(
part_size,
read_costs,
deferred_stripe_handles,
deferred_reopeners,
)
.await;
// Dropping `write_half` on return signals EOF to the reader half.
@@ -3236,6 +3248,7 @@ mod metadata_cache_tests {
mod tests {
use super::*;
use crate::erasure::coding::BitrotWriter;
use serial_test::serial;
use std::io::{Cursor, ErrorKind, IoSlice};
use std::sync::{
Arc,
@@ -3246,6 +3259,15 @@ mod tests {
const CODEC_STREAMING_TEST_BUCKET: &str = "bucket";
const CODEC_STREAMING_TEST_OBJECT: &str = "object";
#[test]
#[serial]
fn multipart_reader_setup_prefetch_is_disabled_only_for_copy_sources() {
temp_env::with_var(ENV_RUSTFS_GET_MULTIPART_READER_SETUP_PREFETCH, Some("true"), || {
assert!(multipart_reader_setup_prefetch_enabled(GetObjectReadPolicy::Default));
assert!(!multipart_reader_setup_prefetch_enabled(GetObjectReadPolicy::CopySource));
});
}
#[tokio::test]
async fn downstream_writer_marks_closed_duplex_reader_as_downstream_close() {
let (reader, inner) = tokio::io::duplex(64);
@@ -5481,6 +5503,7 @@ mod tests {
erasure,
setup.readers,
setup.deferred_stripe_handles,
Vec::new(),
None,
0,
data.len(),
@@ -5531,6 +5554,7 @@ mod tests {
erasure,
setup.readers,
setup.deferred_stripe_handles,
Vec::new(),
None,
0,
part2_len,
@@ -5571,6 +5595,7 @@ mod tests {
erasure,
setup.readers,
setup.deferred_stripe_handles,
Vec::new(),
None,
0,
data.len(),
@@ -6007,6 +6032,49 @@ mod tests {
);
}
#[tokio::test]
#[serial]
async fn codec_streaming_reader_gate_keeps_copy_source_demand_bound_for_all_classes() {
temp_env::async_with_vars(
[
(ENV_RUSTFS_GET_CODEC_STREAMING_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_ROLLOUT, Some("benchmark")),
(ENV_RUSTFS_GET_CODEC_STREAMING_BODY_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_HEADER_COMPAT_CONFIRMED, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MULTIPART_ENABLE, Some("true")),
(ENV_RUSTFS_GET_CODEC_STREAMING_MIN_SIZE, Some("1")),
],
async {
let fi = codec_streaming_test_fileinfo(1024, 2);
let object_info = codec_streaming_test_object_info(&fi);
let normal = codec_streaming_reader_gate_for_test(&None, &object_info, &fi, true);
assert_eq!(normal.decision, GetCodecStreamingDecision::Use);
let plain_fi = codec_streaming_test_fileinfo(1024, 1);
let plain_object_info = codec_streaming_test_object_info(&plain_fi);
let normal_plain = codec_streaming_reader_gate_for_test(&None, &plain_object_info, &plain_fi, true);
assert_eq!(normal_plain.object_class, GetCodecStreamingObjectClass::PlainSinglePart);
assert_eq!(normal_plain.decision, GetCodecStreamingDecision::Use);
let copy = with_get_object_read_policy(GetObjectReadPolicy::CopySource, async {
let multipart = codec_streaming_reader_gate_for_test(&None, &object_info, &fi, true);
let plain = codec_streaming_reader_gate_for_test(&None, &plain_object_info, &plain_fi, true);
(multipart, plain)
})
.await;
assert_eq!(
copy.0.decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::CopySourceDemandBound)
);
assert_eq!(
copy.1.decision,
GetCodecStreamingDecision::Fallback(GetCodecStreamingFallbackReason::CopySourceDemandBound)
);
},
)
.await;
}
#[test]
fn codec_streaming_reader_gate_keeps_multipart_default_off() {
temp_env::with_vars(
@@ -6108,6 +6176,10 @@ mod tests {
assert_eq!(GetCodecStreamingFallbackReason::InvalidMinSize.as_str(), "invalid_min_size");
assert_eq!(GetCodecStreamingFallbackReason::ReadQuorumNotSafe.as_str(), "read_quorum_not_safe");
assert_eq!(GetCodecStreamingFallbackReason::MultipartPartLimit.as_str(), "multipart_part_limit");
assert_eq!(
GetCodecStreamingFallbackReason::CopySourceDemandBound.as_str(),
"copy_source_demand_bound"
);
assert_eq!(GetCodecStreamingObjectClass::PlainSinglePart.as_str(), "plain_single_part");
assert_eq!(GetCodecStreamingObjectClass::Range.as_str(), "range");
assert_eq!(GetCodecStreamingObjectClass::Encrypted.as_str(), "encrypted");
+27
View File
@@ -2468,6 +2468,33 @@ impl ECStore {
Self::resolve_decommission_tiered_object_result(result, bucket, &object)
}
/// Open a source reader for a server-side copy.
///
/// Copy consumers hold the source reader while a destination write can
/// apply backpressure. Keep that read contract explicit at the storage
/// boundary so the lower-level legacy multipart pipeline can suppress its
/// speculative next-part setup without changing the public `ObjectIO`
/// trait or `ObjectOptions` layout.
pub async fn get_object_reader_for_copy(
&self,
bucket: &str,
object: &str,
range: Option<HTTPRangeSpec>,
h: HeaderMap,
opts: &ObjectOptions,
) -> Result<(GetObjectReader, tokio_util::sync::CancellationToken)> {
let cancellation = tokio_util::sync::CancellationToken::new();
let reader = crate::set_disk::with_get_object_read_cancellation(
cancellation.clone(),
crate::set_disk::with_get_object_read_policy(
crate::set_disk::GetObjectReadPolicy::CopySource,
self.handle_get_object_reader(bucket, object, range, h, opts),
),
)
.await?;
Ok((reader, cancellation))
}
#[instrument(level = "debug", skip(self, h))]
#[hotpath::measure(impl_type = "ECStore")]
pub(super) async fn handle_get_object_reader(
+5 -3
View File
@@ -33,7 +33,9 @@ use super::storage_api::multipart_usecase::contract::http::HTTPPreconditions;
use super::storage_api::multipart_usecase::contract::multipart::{
CompletePart, MAX_MULTIPART_PART_NUMBER, MultipartOperations as _, MultipartUploadResult,
};
use super::storage_api::multipart_usecase::contract::object::{ObjectIO as _, ObjectOperations as _};
#[cfg(test)]
use super::storage_api::multipart_usecase::contract::object::ObjectIO as _;
use super::storage_api::multipart_usecase::contract::object::ObjectOperations as _;
use super::storage_api::multipart_usecase::contract::range::HTTPRangeSpec;
use super::storage_api::multipart_usecase::data_usage::{
quota_object_size, record_bucket_object_version_write_memory, record_bucket_object_write_memory,
@@ -1443,8 +1445,8 @@ impl DefaultMultipartUsecase {
.into());
}
let src_reader = store
.get_object_reader(&src_bucket, &src_key, rs.clone(), h, &get_opts)
let (src_reader, _source_cancellation) = store
.get_object_reader_for_copy(&src_bucket, &src_key, rs.clone(), h, &get_opts)
.await
.map_err(map_get_object_reader_error)?;
+9 -2
View File
@@ -7948,11 +7948,18 @@ impl DefaultObjectUsecase {
.into());
}
let gr = store
.get_object_reader(&src_bucket, &src_key, None, h, &src_get_opts)
let (gr, source_cancellation) = store
.get_object_reader_for_copy(&src_bucket, &src_key, None, h, &src_get_opts)
.await
.map_err(map_get_object_reader_error)?;
// The commit owner is intentionally detached so SetDisk can finish
// its rename/cleanup and post-commit publication if the HTTP caller
// goes away. Keep a request-owned guard for the source producer:
// cancellation drops the source read promptly, while the detached
// commit task retains the guards it needs to complete safely.
let _source_cancellation_guard = source_cancellation.clone().drop_guard();
let mut src_info = gr.object_info.clone();
// A copy reads the source plaintext, so it needs the source key's decrypt permission
+3 -1
View File
@@ -1185,7 +1185,9 @@ pub(crate) mod multipart_usecase {
}
pub(crate) mod object {
pub(crate) use super::super::super::storage_contracts::{ObjectIO, ObjectOperations};
#[cfg(test)]
pub(crate) use super::super::super::storage_contracts::ObjectIO;
pub(crate) use super::super::super::storage_contracts::ObjectOperations;
}
pub(crate) mod range {
+1
View File
@@ -710,6 +710,7 @@ Static gate reasons:
- multipart_part_limit
- invalid_min_size
- read_quorum_not_safe
- copy_source_demand_bound
Relevant focused proof points:
- set_disk::read::tests::codec_streaming_reader_gate_defaults_to_disabled