diff --git a/crates/ecstore/src/cluster/rpc/internode_data_transport.rs b/crates/ecstore/src/cluster/rpc/internode_data_transport.rs index 9c33d191a..6f46fead3 100644 --- a/crates/ecstore/src/cluster/rpc/internode_data_transport.rs +++ b/crates/ecstore/src/cluster/rpc/internode_data_transport.rs @@ -233,11 +233,17 @@ pub struct NsScannerCapabilityRequest { #[async_trait] pub trait InternodeDataTransport: Send + Sync + std::fmt::Debug { async fn open_read(&self, request: ReadStreamRequest) -> Result; + async fn open_read_fresh(&self, request: ReadStreamRequest) -> Result { + self.open_read(request).await + } /// Opens an owned-chunk stream when this transport can retain receive-buffer /// ownership. `None` preserves the established `open_read` fallback. async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result> { Ok(None) } + async fn open_read_chunks_fresh(&self, request: ReadStreamRequest) -> Result> { + self.open_read_chunks(request).await + } async fn open_write(&self, request: WriteStreamRequest) -> Result; async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result; async fn open_ns_scanner(&self, _request: NsScannerStreamRequest) -> Result { @@ -269,6 +275,15 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport { )) } + async fn open_read_fresh(&self, request: ReadStreamRequest) -> Result { + let url = build_read_file_stream_url(&request); + let mut headers = json_headers(); + build_auth_headers(&url, &Method::GET, &mut headers)?; + Ok(Box::new( + HttpReader::new_fresh_connection_with_stall_timeout(url, Method::GET, headers, None, request.stall_timeout).await?, + )) + } + async fn open_read_chunks(&self, request: ReadStreamRequest) -> Result> { let url = build_read_file_stream_url(&request); let mut headers = json_headers(); @@ -278,6 +293,16 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport { ))) } + async fn open_read_chunks_fresh(&self, request: ReadStreamRequest) -> Result> { + let url = build_read_file_stream_url(&request); + let mut headers = json_headers(); + build_auth_headers(&url, &Method::GET, &mut headers)?; + Ok(Some(Box::new( + HttpChunkReader::new_fresh_connection_with_stall_timeout(url, Method::GET, headers, None, request.stall_timeout) + .await?, + ))) + } + async fn open_write(&self, request: WriteStreamRequest) -> Result { let server_epoch = self.put_file_auth_capability(&request.endpoint).await?; let nonce = server_epoch.map(|_| Uuid::new_v4()); diff --git a/crates/ecstore/src/cluster/rpc/remote_disk.rs b/crates/ecstore/src/cluster/rpc/remote_disk.rs index 13e45fbd3..3ae2c8b01 100644 --- a/crates/ecstore/src/cluster/rpc/remote_disk.rs +++ b/crates/ecstore/src/cluster/rpc/remote_disk.rs @@ -55,18 +55,22 @@ use rustfs_protos::proto_gen::node_service::{ }; use serde::{Serialize, de::DeserializeOwned}; use std::{ + future::Future, io::Cursor, path::PathBuf, + pin::Pin, sync::{ Arc, atomic::{AtomicU32, Ordering}, }, + task::{Context, Poll}, time::Duration, }; use tokio::time; use tokio::{ - io::{self, AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt}, + io::{self, AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, ReadBuf}, net::TcpStream, + task::{JoinError, JoinHandle}, time::timeout, }; use tokio_util::sync::CancellationToken; @@ -84,6 +88,7 @@ const REMOTE_DISK_OPEN_WRITE_MAX_ATTEMPTS: usize = 2; const REMOTE_DISK_OPEN_WRITE_RETRY_BACKOFF: Duration = Duration::from_millis(20); const REMOTE_DISK_OPEN_READ_MAX_ATTEMPTS: usize = 2; const REMOTE_DISK_OPEN_READ_RETRY_BACKOFF: Duration = Duration::from_millis(20); +const REMOTE_READ_TIMEOUT_PARTS: u32 = 3; const NS_SCANNER_CAPABILITY_PROBE_TIMEOUT: Duration = Duration::from_secs(5); /// Base backoff for idempotent read-only RPC retries (grpc-optimization P3-3); doubles per attempt. const REMOTE_DISK_READ_RETRY_BASE_BACKOFF: Duration = Duration::from_millis(50); @@ -214,6 +219,415 @@ where } } +fn is_retryable_remote_body_error(error: &io::Error) -> bool { + if error + .get_ref() + .and_then(|source| source.downcast_ref::()) + .is_some() + { + return true; + } + + matches!( + error.kind(), + io::ErrorKind::ConnectionReset + | io::ErrorKind::BrokenPipe + | io::ErrorKind::ConnectionAborted + | io::ErrorKind::UnexpectedEof + ) +} + +fn resumed_read_request(request: &ReadStreamRequest, emitted: usize) -> io::Result { + let offset = request + .offset + .checked_add(emitted) + .ok_or_else(|| io::Error::other("remote read resume offset overflow"))?; + let length = if request.length == 0 { + 0 + } else { + request + .length + .checked_sub(emitted) + .ok_or_else(|| io::Error::other("remote read resume offset exceeds requested length"))? + }; + Ok(ReadStreamRequest { + offset, + length, + ..request.clone() + }) +} + +#[derive(Clone, Copy)] +struct RemoteReadTimeouts { + body_stall: Option, + initial_read: Option, + recovery: Option, +} + +fn remote_read_timeouts(read_timeout: Duration) -> RemoteReadTimeouts { + let Some(recovery) = read_timeout + .checked_div(REMOTE_READ_TIMEOUT_PARTS) + .filter(|timeout| !timeout.is_zero()) + else { + return RemoteReadTimeouts { + body_stall: None, + initial_read: None, + recovery: None, + }; + }; + RemoteReadTimeouts { + body_stall: Some(recovery), + initial_read: Some(read_timeout.saturating_sub(recovery)), + recovery: Some(recovery), + } +} + +async fn with_remote_read_recovery_timeout(recovery_timeout: Option, future: F) -> Result +where + F: Future>, +{ + match recovery_timeout { + Some(recovery_timeout) => match time::timeout(recovery_timeout, future).await { + Ok(result) => result, + Err(_) => Err(DiskError::Timeout), + }, + None => future.await, + } +} + +struct AbortOnDropTask(JoinHandle); + +impl AbortOnDropTask { + fn new(handle: JoinHandle) -> Self { + Self(handle) + } +} + +impl Future for AbortOnDropTask +where + T: Send + 'static, +{ + type Output = std::result::Result; + + fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + Pin::new(&mut self.get_mut().0).poll(cx) + } +} + +impl Drop for AbortOnDropTask { + fn drop(&mut self) { + self.0.abort(); + } +} + +fn retry_cutoff_elapsed( + initial_read_timeout: &mut Option, + cutoff: &mut Option>>, + cx: &mut Context<'_>, +) -> bool { + if cutoff.is_none() + && let Some(timeout) = initial_read_timeout.take() + { + *cutoff = Some(Box::pin(time::sleep(timeout))); + } + cutoff.as_mut().is_some_and(|cutoff| cutoff.as_mut().poll(cx).is_ready()) +} + +type ReadResumeFuture = AbortOnDropTask>; + +struct RetryingRemoteReader { + reader: Option, + transport: Arc, + request: ReadStreamRequest, + emitted: usize, + retried: bool, + initial_read_timeout: Option, + retry_cutoff: Option>>, + recovery_timeout: Option, + resume: Option, +} + +impl RetryingRemoteReader { + fn new_with_timeouts( + reader: FileReader, + transport: Arc, + request: ReadStreamRequest, + initial_read_timeout: Option, + recovery_timeout: Option, + ) -> Self { + Self { + reader: Some(reader), + transport, + request, + emitted: 0, + retried: false, + initial_read_timeout, + retry_cutoff: None, + recovery_timeout, + resume: None, + } + } + + fn start_resume(&mut self) -> io::Result<()> { + if self.request.length != 0 && self.emitted >= self.request.length { + self.reader = None; + return Ok(()); + } + let request = resumed_read_request(&self.request, self.emitted)?; + let recovery_timeout = self.recovery_timeout; + let transport = Arc::clone(&self.transport); + self.resume = Some(AbortOnDropTask::new(tokio::spawn(async move { + with_remote_read_recovery_timeout(recovery_timeout, transport.open_read_fresh(request)).await + }))); + Ok(()) + } + + fn retry_cutoff_elapsed(&mut self, cx: &mut Context<'_>) -> bool { + retry_cutoff_elapsed(&mut self.initial_read_timeout, &mut self.retry_cutoff, cx) + } +} + +impl AsyncRead for RetryingRemoteReader { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + 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() { + match Pin::new(resume).poll(cx) { + Poll::Pending => true, + Poll::Ready(Ok(Ok(reader))) => { + self.resume = None; + self.reader = Some(reader); + false + } + Poll::Ready(Ok(Err(error))) => { + self.resume = None; + if self.reader.is_none() { + return Poll::Ready(Err(io::Error::other(error))); + } + continue; + } + Poll::Ready(Err(error)) => { + self.resume = None; + if self.reader.is_none() { + return Poll::Ready(Err(io::Error::other(error))); + } + continue; + } + } + } else { + false + }; + + if !self.retried && self.retry_cutoff_elapsed(cx) { + self.retried = true; + if let Err(resume_error) = self.start_resume() { + return Poll::Ready(Err(resume_error)); + } + continue; + } + + let Some(reader) = self.reader.as_mut() else { + if resume_pending { + return Poll::Pending; + } + return Poll::Ready(Ok(())); + }; + let before = buf.filled().len(); + match Pin::new(reader).poll_read(cx, buf) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Ok(())) => { + let produced = buf.filled().len() - before; + self.emitted = match self.emitted.checked_add(produced) { + Some(emitted) => emitted, + None => return Poll::Ready(Err(io::Error::other("remote read emitted byte count overflow"))), + }; + if resume_pending { + if produced == 0 && (self.request.length == 0 || self.emitted >= self.request.length) { + self.resume = None; + } else if produced == 0 { + self.reader = None; + continue; + } else { + self.resume = None; + } + } + return Poll::Ready(Ok(())); + } + Poll::Ready(Err(error)) if !self.retried && is_retryable_remote_body_error(&error) => { + self.retried = true; + self.reader = None; + if let Err(resume_error) = self.start_resume() { + return Poll::Ready(Err(resume_error)); + } + continue; + } + Poll::Ready(Err(error)) if resume_pending && is_retryable_remote_body_error(&error) => { + self.reader = None; + continue; + } + Poll::Ready(Err(error)) => return Poll::Ready(Err(error)), + } + } + } +} + +type ChunkResumeFuture = AbortOnDropTask>>; + +struct RetryingRemoteChunkReader { + reader: Option, + transport: Arc, + request: ReadStreamRequest, + emitted: usize, + retried: bool, + initial_read_timeout: Option, + retry_cutoff: Option>>, + recovery_timeout: Option, + resume: Option, +} + +impl RetryingRemoteChunkReader { + fn new_with_timeouts( + reader: rustfs_rio::ChunkReaderBox, + transport: Arc, + request: ReadStreamRequest, + initial_read_timeout: Option, + recovery_timeout: Option, + ) -> Self { + Self { + reader: Some(reader), + transport, + request, + emitted: 0, + retried: false, + initial_read_timeout, + retry_cutoff: None, + recovery_timeout, + resume: None, + } + } + + fn start_resume(&mut self) -> io::Result<()> { + if self.request.length != 0 && self.emitted >= self.request.length { + self.reader = None; + return Ok(()); + } + let request = resumed_read_request(&self.request, self.emitted)?; + let recovery_timeout = self.recovery_timeout; + let transport = Arc::clone(&self.transport); + self.resume = Some(AbortOnDropTask::new(tokio::spawn(async move { + with_remote_read_recovery_timeout(recovery_timeout, transport.open_read_chunks_fresh(request)).await + }))); + Ok(()) + } + + fn retry_cutoff_elapsed(&mut self, cx: &mut Context<'_>) -> bool { + retry_cutoff_elapsed(&mut self.initial_read_timeout, &mut self.retry_cutoff, cx) + } +} + +impl AsyncRead for RetryingRemoteChunkReader { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + if buf.remaining() == 0 { + return Poll::Ready(Ok(())); + } + match rustfs_rio::ChunkReader::poll_read_chunk(self.as_mut(), cx, buf.remaining()) { + Poll::Ready(Ok(Some(chunk))) => { + buf.put_slice(&chunk); + Poll::Ready(Ok(())) + } + Poll::Ready(Ok(None)) => Poll::Ready(Ok(())), + Poll::Ready(Err(error)) => Poll::Ready(Err(error)), + Poll::Pending => Poll::Pending, + } + } +} + +impl rustfs_rio::ChunkReader for RetryingRemoteChunkReader { + fn poll_read_chunk(mut self: Pin<&mut Self>, cx: &mut Context<'_>, max: usize) -> Poll>> { + loop { + let resume_pending = if let Some(resume) = self.resume.as_mut() { + match Pin::new(resume).poll(cx) { + Poll::Pending => true, + Poll::Ready(Ok(Ok(Some(reader)))) => { + self.resume = None; + self.reader = Some(reader); + false + } + 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"))); + } + continue; + } + Poll::Ready(Ok(Err(error))) => { + self.resume = None; + if self.reader.is_none() { + return Poll::Ready(Err(io::Error::other(error))); + } + continue; + } + Poll::Ready(Err(error)) => { + self.resume = None; + if self.reader.is_none() { + return Poll::Ready(Err(io::Error::other(error))); + } + continue; + } + } + } else { + false + }; + + if !self.retried && self.retry_cutoff_elapsed(cx) { + self.retried = true; + if let Err(resume_error) = self.start_resume() { + return Poll::Ready(Err(resume_error)); + } + continue; + } + + let Some(reader) = self.reader.as_mut() else { + if resume_pending { + return Poll::Pending; + } + return Poll::Ready(Ok(None)); + }; + match rustfs_rio::ChunkReader::poll_read_chunk(Pin::new(reader.as_mut()), cx, max) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Ok(Some(chunk))) => { + self.emitted = match self.emitted.checked_add(chunk.len()) { + Some(emitted) => emitted, + None => return Poll::Ready(Err(io::Error::other("remote read emitted byte count overflow"))), + }; + if resume_pending { + self.resume = None; + } + return Poll::Ready(Ok(Some(chunk))); + } + Poll::Ready(Ok(None)) if resume_pending => { + self.reader = None; + continue; + } + Poll::Ready(Ok(None)) => return Poll::Ready(Ok(None)), + Poll::Ready(Err(error)) if !self.retried && is_retryable_remote_body_error(&error) => { + self.retried = true; + self.reader = None; + if let Err(resume_error) = self.start_resume() { + return Poll::Ready(Err(resume_error)); + } + continue; + } + Poll::Ready(Err(error)) if resume_pending && is_retryable_remote_body_error(&error) => { + self.reader = None; + continue; + } + Poll::Ready(Err(error)) => return Poll::Ready(Err(error)), + } + } + } +} + #[derive(Debug)] pub struct RemoteDisk { pub id: Mutex>, @@ -2483,17 +2897,24 @@ impl DiskAPI for RemoteDisk { return Err(DiskError::FaultyDisk); } let disk = self.disk_ref().await; - let stall_timeout = get_object_disk_read_timeout(); - self.open_read_with_retry(ReadStreamRequest { + let timeouts = remote_read_timeouts(get_object_disk_read_timeout()); + let request = ReadStreamRequest { endpoint: self.endpoint.grid_host(), disk, volume: volume.to_string(), path: path.to_string(), offset, length, - stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout), - }) - .await + stall_timeout: timeouts.body_stall, + }; + let reader = self.open_read_with_retry(request.clone()).await?; + Ok(Box::new(RetryingRemoteReader::new_with_timeouts( + reader, + Arc::clone(&self.data_transport), + request, + timeouts.initial_read, + timeouts.recovery, + ))) } async fn read_file_stream_chunks( @@ -2507,17 +2928,26 @@ impl DiskAPI for RemoteDisk { return Err(DiskError::FaultyDisk); } let disk = self.disk_ref().await; - let stall_timeout = get_object_disk_read_timeout(); - self.open_read_chunks_with_retry(ReadStreamRequest { + let timeouts = remote_read_timeouts(get_object_disk_read_timeout()); + let request = ReadStreamRequest { endpoint: self.endpoint.grid_host(), disk, volume: volume.to_string(), path: path.to_string(), offset, length, - stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout), - }) - .await + stall_timeout: timeouts.body_stall, + }; + let reader = self.open_read_chunks_with_retry(request.clone()).await?; + Ok(reader.map(|reader| { + Box::new(RetryingRemoteChunkReader::new_with_timeouts( + reader, + Arc::clone(&self.data_transport), + request, + timeouts.initial_read, + timeouts.recovery, + )) as rustfs_rio::ChunkReaderBox + })) } /// Buffered read for remote disks. @@ -3115,12 +3545,14 @@ impl DiskAPI for RemoteDisk { mod tests { use super::*; use crate::cluster::rpc::internode_data_transport::{InternodeDataTransportCapabilities, TcpHttpInternodeDataTransport}; + use crate::erasure::coding::{BitrotReader, Erasure, decode::ParallelReader}; + use crate::io_support::bitrot::ShardReader; use crate::runtime::sources as runtime_sources; use serde_json::Value; use serial_test::serial; use std::io::{self as std_io, Write}; use std::pin::Pin; - use std::sync::{Arc, Mutex, Mutex as StdMutex, Once}; + use std::sync::{Arc, Mutex, Mutex as StdMutex, Once, atomic::AtomicUsize}; use std::task::{Context, Poll}; use tokio::io::{ReadBuf, duplex}; use tokio::net::TcpListener; @@ -4138,6 +4570,735 @@ mod tests { } } + #[derive(Debug, Clone)] + enum ResumeReadStep { + PartialThenReset(Vec), + Data(Vec), + } + + #[derive(Debug, Default)] + struct ResumeTransport { + read_steps: Mutex>, + chunk_steps: Mutex>, + read_requests: Mutex>, + chunk_requests: Mutex>, + fresh_read_requests: Mutex>, + fresh_chunk_requests: Mutex>, + } + + impl ResumeTransport { + fn with_read_steps(read_steps: Vec) -> Self { + Self { + read_steps: Mutex::new(read_steps), + ..Self::default() + } + } + + fn with_chunk_steps(chunk_steps: Vec) -> Self { + Self { + chunk_steps: Mutex::new(chunk_steps), + ..Self::default() + } + } + } + + #[derive(Debug)] + struct ChunkPartialThenErrorReader { + data: Option, + error: Option, + } + + impl rustfs_rio::ChunkReader for ChunkPartialThenErrorReader { + fn poll_read_chunk(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, max: usize) -> Poll>> { + if let Some(mut data) = self.data.take() { + let take = data.len().min(max); + let chunk = data.split_to(take); + if !data.is_empty() { + self.data = Some(data); + } + return Poll::Ready(Ok(Some(chunk))); + } + if let Some(error) = self.error.take() { + return Poll::Ready(Err(error)); + } + Poll::Ready(Ok(None)) + } + } + + impl AsyncRead for ChunkPartialThenErrorReader { + fn poll_read(self: Pin<&mut Self>, _cx: &mut Context<'_>, _buf: &mut ReadBuf<'_>) -> Poll> { + Poll::Ready(Err(io::Error::other("chunk reader must use chunk handoff"))) + } + } + + #[derive(Debug)] + struct BodyStallTestReader { + data: Option, + next_data: Option<(Bytes, Pin>)>, + initial_delay: Option>>, + timeout: Duration, + stall_timer: Option>>, + } + + impl BodyStallTestReader { + fn new(data: Bytes, initial_delay: Duration, timeout: Option) -> Self { + let timeout = timeout.expect("parallel resume test requires a body stall timeout"); + Self { + data: Some(data), + next_data: None, + initial_delay: Some(Box::pin(time::sleep(initial_delay))), + timeout, + stall_timer: None, + } + } + + fn with_next_data( + data: Bytes, + initial_delay: Duration, + next_data: Bytes, + next_delay: Duration, + timeout: Option, + ) -> Self { + let mut reader = Self::new(data, initial_delay, timeout); + reader.next_data = Some((next_data, Box::pin(time::sleep(next_delay)))); + reader + } + + fn poll_chunk(&mut self, cx: &mut Context<'_>, max: usize) -> Poll>> { + if let Some(delay) = self.initial_delay.as_mut() { + if delay.as_mut().poll(cx).is_pending() { + return Poll::Pending; + } + self.initial_delay = None; + } + if let Some(mut data) = self.data.take() { + let chunk = data.split_to(data.len().min(max)); + if !data.is_empty() { + self.data = Some(data); + } + return Poll::Ready(Ok(Some(chunk))); + } + if let Some((_, delay)) = self.next_data.as_mut() + && delay.as_mut().poll(cx).is_pending() + { + return Poll::Pending; + } + if let Some((mut data, _)) = self.next_data.take() { + let chunk = data.split_to(data.len().min(max)); + if !data.is_empty() { + self.next_data = Some((data, Box::pin(time::sleep(Duration::ZERO)))); + } + return Poll::Ready(Ok(Some(chunk))); + } + let timer = self.stall_timer.get_or_insert_with(|| Box::pin(time::sleep(self.timeout))); + match timer.as_mut().poll(cx) { + Poll::Pending => Poll::Pending, + Poll::Ready(()) => Poll::Ready(Err(io::Error::new( + std_io::ErrorKind::TimedOut, + rustfs_rio::BodyStalled { timeout: self.timeout }, + ))), + } + } + } + + impl AsyncRead for BodyStallTestReader { + fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll> { + match self.poll_chunk(cx, buf.remaining()) { + Poll::Ready(Ok(Some(chunk))) => { + buf.put_slice(&chunk); + Poll::Ready(Ok(())) + } + Poll::Ready(Ok(None)) => Poll::Ready(Ok(())), + Poll::Ready(Err(error)) => Poll::Ready(Err(error)), + Poll::Pending => Poll::Pending, + } + } + } + + impl rustfs_rio::ChunkReader for BodyStallTestReader { + fn poll_read_chunk(mut self: Pin<&mut Self>, cx: &mut Context<'_>, max: usize) -> Poll>> { + self.poll_chunk(cx, max) + } + } + + struct CountOnDrop(Arc); + + impl Drop for CountOnDrop { + fn drop(&mut self) { + self.0.fetch_add(1, Ordering::Relaxed); + } + } + + #[derive(Debug)] + struct ParallelResumeTransport { + initial_data: Bytes, + resumed_data: Bytes, + initial_delay: Duration, + fresh_delay: Duration, + fresh_read_requests: Mutex>, + fresh_chunk_requests: Mutex>, + initial_next_data: Option<(Bytes, Duration)>, + } + + impl ParallelResumeTransport { + fn new(initial_delay: Duration, fresh_delay: Duration) -> Self { + Self { + initial_data: Bytes::from_static(b"da"), + resumed_data: Bytes::from_static(b"ta"), + initial_delay, + fresh_delay, + fresh_read_requests: Mutex::new(Vec::new()), + fresh_chunk_requests: Mutex::new(Vec::new()), + initial_next_data: None, + } + } + + fn with_initial_next_data(initial_delay: Duration, next_delay: Duration, fresh_delay: Duration) -> Self { + let mut transport = Self::new(initial_delay, fresh_delay); + transport.initial_next_data = Some((Bytes::from_static(b"ta"), next_delay)); + transport + } + } + + #[async_trait::async_trait] + impl InternodeDataTransport for ParallelResumeTransport { + async fn open_read(&self, request: ReadStreamRequest) -> Result { + let reader = match self.initial_next_data.as_ref() { + Some((next_data, next_delay)) => BodyStallTestReader::with_next_data( + self.initial_data.clone(), + self.initial_delay, + next_data.clone(), + *next_delay, + request.stall_timeout, + ), + None => BodyStallTestReader::new(self.initial_data.clone(), self.initial_delay, request.stall_timeout), + }; + Ok(Box::new(reader)) + } + + async fn open_read_fresh(&self, request: ReadStreamRequest) -> Result { + self.fresh_read_requests + .lock() + .expect("fresh read request lock should not be poisoned") + .push(request); + time::sleep(self.fresh_delay).await; + Ok(Box::new(Cursor::new(self.resumed_data.clone()))) + } + + async fn open_read_chunks(&self, request: ReadStreamRequest) -> Result> { + let reader = match self.initial_next_data.as_ref() { + Some((next_data, next_delay)) => BodyStallTestReader::with_next_data( + self.initial_data.clone(), + self.initial_delay, + next_data.clone(), + *next_delay, + request.stall_timeout, + ), + None => BodyStallTestReader::new(self.initial_data.clone(), self.initial_delay, request.stall_timeout), + }; + Ok(Some(Box::new(reader))) + } + + async fn open_read_chunks_fresh(&self, request: ReadStreamRequest) -> Result> { + self.fresh_chunk_requests + .lock() + .expect("fresh chunk request lock should not be poisoned") + .push(request); + time::sleep(self.fresh_delay).await; + Ok(Some(Box::new(ChunkPartialThenErrorReader { + data: Some(self.resumed_data.clone()), + error: None, + }))) + } + + async fn open_write(&self, _request: WriteStreamRequest) -> Result { + panic!("open_write should not be used in parallel resume tests"); + } + + async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result { + panic!("open_walk_dir should not be used in parallel resume tests"); + } + + fn name(&self) -> &'static str { + "parallel-resume-test" + } + + fn capabilities(&self) -> InternodeDataTransportCapabilities { + InternodeDataTransportCapabilities::tcp_http() + } + } + + #[derive(Debug, Default)] + struct PendingFreshOpenTransport { + fresh_read_drops: Arc, + fresh_chunk_drops: Arc, + } + + #[async_trait::async_trait] + impl InternodeDataTransport for PendingFreshOpenTransport { + async fn open_read(&self, _request: ReadStreamRequest) -> Result { + 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 { + let _drop = CountOnDrop(Arc::clone(&self.fresh_read_drops)); + std::future::pending().await + } + + async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result> { + Ok(Some(Box::new(ChunkPartialThenErrorReader { + data: None, + error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")), + }))) + } + + async fn open_read_chunks_fresh(&self, _request: ReadStreamRequest) -> Result> { + let _drop = CountOnDrop(Arc::clone(&self.fresh_chunk_drops)); + std::future::pending().await + } + + async fn open_write(&self, _request: WriteStreamRequest) -> Result { + panic!("open_write should not be used in fresh open cancellation tests"); + } + + async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result { + panic!("open_walk_dir should not be used in fresh open cancellation tests"); + } + + fn name(&self) -> &'static str { + "pending-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 { + cursor: Cursor::new(data), + error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")), + }), + ResumeReadStep::Data(data) => Box::new(Cursor::new(data)), + } + } + + fn resume_step_chunk_reader(step: ResumeReadStep) -> rustfs_rio::ChunkReaderBox { + match step { + ResumeReadStep::PartialThenReset(data) => Box::new(ChunkPartialThenErrorReader { + data: Some(Bytes::from(data)), + error: Some(io::Error::new(std_io::ErrorKind::ConnectionReset, "stream reset")), + }), + ResumeReadStep::Data(data) => Box::new(ChunkPartialThenErrorReader { + data: Some(Bytes::from(data)), + error: None, + }), + } + } + + #[async_trait::async_trait] + impl InternodeDataTransport for ResumeTransport { + async fn open_read(&self, request: ReadStreamRequest) -> Result { + self.read_requests + .lock() + .expect("read request lock should not be poisoned") + .push(request); + let step = self + .read_steps + .lock() + .expect("read steps lock should not be poisoned") + .remove(0); + Ok(resume_step_reader(step)) + } + + async fn open_read_fresh(&self, request: ReadStreamRequest) -> Result { + self.fresh_read_requests + .lock() + .expect("fresh read request lock should not be poisoned") + .push(request.clone()); + self.open_read(request).await + } + + async fn open_read_chunks(&self, request: ReadStreamRequest) -> Result> { + self.chunk_requests + .lock() + .expect("chunk request lock should not be poisoned") + .push(request); + let step = self + .chunk_steps + .lock() + .expect("chunk steps lock should not be poisoned") + .remove(0); + Ok(Some(resume_step_chunk_reader(step))) + } + + async fn open_read_chunks_fresh(&self, request: ReadStreamRequest) -> Result> { + self.fresh_chunk_requests + .lock() + .expect("fresh chunk request lock should not be poisoned") + .push(request.clone()); + self.open_read_chunks(request).await + } + + async fn open_write(&self, _request: WriteStreamRequest) -> Result { + panic!("open_write should not be used in remote read resume tests"); + } + + async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result { + panic!("open_walk_dir should not be used in remote read resume tests"); + } + + fn name(&self) -> &'static str { + "resume-test" + } + + fn capabilities(&self) -> InternodeDataTransportCapabilities { + InternodeDataTransportCapabilities::tcp_http() + } + } + + fn resume_request(length: usize) -> ReadStreamRequest { + ReadStreamRequest { + endpoint: "http://remote".to_string(), + disk: "disk".to_string(), + volume: "volume".to_string(), + path: "path".to_string(), + offset: 7, + length, + stall_timeout: None, + } + } + + #[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())])); + let request = resume_request(10); + let reader = resume_step_reader(ResumeReadStep::PartialThenReset(b"0123".to_vec())); + let mut reader = RetryingRemoteReader::new_with_timeouts(reader, transport.clone(), request, None, None); + let mut output = Vec::new(); + reader + .read_to_end(&mut output) + .await + .expect("one body reset should be resumed"); + + assert_eq!(output, b"0123456789"); + let requests = transport + .read_requests + .lock() + .expect("read request lock should not be poisoned"); + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].offset, 11); + assert_eq!(requests[0].length, 6); + 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_resumes_from_emitted_bytes_without_duplicates() { + let transport = Arc::new(ResumeTransport::with_chunk_steps(vec![ResumeReadStep::Data(b"456789".to_vec())])); + let request = resume_request(10); + let reader = resume_step_chunk_reader(ResumeReadStep::PartialThenReset(b"0123".to_vec())); + let mut reader = RetryingRemoteChunkReader::new_with_timeouts(reader, transport.clone(), request, None, None); + let mut output = Vec::new(); + reader + .read_to_end(&mut output) + .await + .expect("chunk body reset should be resumed"); + + assert_eq!(output, b"0123456789"); + let requests = transport + .chunk_requests + .lock() + .expect("chunk request lock should not be poisoned"); + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].offset, 11); + assert_eq!(requests[0].length, 6); + assert_eq!( + transport + .fresh_chunk_requests + .lock() + .expect("fresh chunk request lock should not be poisoned") + .len(), + 1 + ); + } + + #[derive(Clone, Copy)] + enum ParallelResumePath { + Regular, + Chunk, + } + + async fn assert_parallel_resume_case( + path: ParallelResumePath, + initial_delay: Duration, + initial_next_delay: Option, + fresh_delay: Duration, + expect_success: bool, + ) { + const DATA: &[u8] = b"data"; + let transport = Arc::new(match initial_next_delay { + Some(next_delay) => ParallelResumeTransport::with_initial_next_data(initial_delay, next_delay, fresh_delay), + None => ParallelResumeTransport::new(initial_delay, fresh_delay), + }); + let remote_disk = new_remote_disk_with_transport(transport.clone()).await; + let erasure = Erasure::new(1, 1, DATA.len()); + let (buffers, errors) = match path { + ParallelResumePath::Regular => { + let reader = remote_disk + .read_file_stream("bucket", "object/part.1", 0, DATA.len()) + .await + .expect("initial remote reader should open"); + let readers = vec![ + Some(BitrotReader::new(reader, DATA.len(), rustfs_utils::HashAlgorithm::None, false)), + None, + ]; + ParallelReader::new_with_metrics_path_and_reconstruction_verification(readers, erasure, 0, DATA.len(), None) + .read() + .await + } + ParallelResumePath::Chunk => { + let reader = remote_disk + .read_file_stream_chunks("bucket", "object/part.1", 0, DATA.len()) + .await + .expect("initial remote chunk reader should open") + .expect("chunk transport should return a reader"); + let readers = vec![ + Some(BitrotReader::new( + ShardReader::Chunked(reader), + DATA.len(), + rustfs_utils::HashAlgorithm::None, + false, + )), + None, + ]; + ParallelReader::new_with_metrics_path_and_reconstruction_verification(readers, erasure, 0, DATA.len(), None) + .read() + .await + } + }; + + if expect_success { + assert_eq!(buffers[0].as_deref(), Some(DATA)); + assert!(errors[0].is_none()); + } else { + assert!(buffers[0].is_none()); + assert!(matches!(errors[0], Some(DiskError::Timeout))); + } + let requests = match path { + ParallelResumePath::Regular => transport + .fresh_read_requests + .lock() + .expect("fresh read request lock should not be poisoned"), + ParallelResumePath::Chunk => transport + .fresh_chunk_requests + .lock() + .expect("fresh chunk request lock should not be poisoned"), + }; + assert_eq!(requests.len(), 1); + assert_eq!(requests[0].offset, 2); + assert_eq!(requests[0].length, 2); + } + + async fn assert_parallel_resume_path(path: ParallelResumePath) { + assert_parallel_resume_case(path, Duration::ZERO, None, Duration::from_millis(100), true).await; + assert_parallel_resume_case(path, Duration::from_millis(500), None, Duration::from_millis(100), true).await; + assert_parallel_resume_case(path, Duration::ZERO, None, Duration::from_millis(500), false).await; + assert_parallel_resume_case( + path, + Duration::from_millis(650), + Some(Duration::from_millis(700)), + Duration::from_millis(500), + true, + ) + .await; + } + + async fn assert_retry_drop_cancels_fresh_open(path: ParallelResumePath) { + let transport = Arc::new(PendingFreshOpenTransport::default()); + let remote_disk = new_remote_disk_with_transport(transport.clone()).await; + let mut output = [0_u8; 1]; + match path { + ParallelResumePath::Regular => { + let mut reader = remote_disk + .read_file_stream("bucket", "object/part.1", 0, 1) + .await + .expect("initial remote reader should open"); + assert!( + time::timeout(Duration::from_millis(20), reader.read(&mut output)) + .await + .is_err() + ); + drop(reader); + } + ParallelResumePath::Chunk => { + let mut reader = remote_disk + .read_file_stream_chunks("bucket", "object/part.1", 0, 1) + .await + .expect("initial remote chunk reader should open") + .expect("chunk transport should return a reader"); + assert!( + time::timeout(Duration::from_millis(20), reader.read(&mut output)) + .await + .is_err() + ); + drop(reader); + } + } + time::timeout(Duration::from_secs(1), async { + loop { + let drops = match path { + ParallelResumePath::Regular => transport.fresh_read_drops.load(Ordering::Relaxed), + ParallelResumePath::Chunk => transport.fresh_chunk_drops.load(Ordering::Relaxed), + }; + if drops == 1 { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("dropping the retrying reader should cancel the pending fresh open"); + } + + #[tokio::test(start_paused = true)] + #[serial] + async fn remote_reader_recovers_body_stall_through_parallel_reader() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_DISK_READ_TIMEOUT, Some("1"))], async { + assert_parallel_resume_path(ParallelResumePath::Regular).await; + }) + .await; + } + + #[tokio::test(start_paused = true)] + #[serial] + async fn remote_chunk_reader_recovers_body_stall_through_parallel_reader() { + temp_env::async_with_vars([(rustfs_config::ENV_OBJECT_DISK_READ_TIMEOUT, Some("1"))], async { + assert_parallel_resume_path(ParallelResumePath::Chunk).await; + }) + .await; + } + + #[tokio::test] + async fn remote_reader_drop_cancels_pending_fresh_open() { + assert_retry_drop_cancels_fresh_open(ParallelResumePath::Regular).await; + } + + #[tokio::test] + async fn remote_chunk_reader_drop_cancels_pending_fresh_open() { + assert_retry_drop_cancels_fresh_open(ParallelResumePath::Chunk).await; + } + + #[tokio::test] + async fn remote_reader_treats_error_after_requested_length_as_eof() { + let transport = Arc::new(ResumeTransport::default()); + let reader = resume_step_reader(ResumeReadStep::PartialThenReset(b"0123".to_vec())); + let mut reader = RetryingRemoteReader::new_with_timeouts(reader, transport.clone(), resume_request(4), None, None); + let mut output = Vec::new(); + + reader + .read_to_end(&mut output) + .await + .expect("error after the requested bytes should not trigger a redundant resume"); + assert_eq!(output, b"0123"); + assert!( + transport + .fresh_read_requests + .lock() + .expect("fresh read request lock should not be poisoned") + .is_empty() + ); + } + + #[tokio::test] + async fn remote_chunk_reader_treats_error_after_requested_length_as_eof() { + let transport = Arc::new(ResumeTransport::default()); + let reader = resume_step_chunk_reader(ResumeReadStep::PartialThenReset(b"0123".to_vec())); + let mut reader = RetryingRemoteChunkReader::new_with_timeouts(reader, transport.clone(), resume_request(4), None, None); + let mut output = Vec::new(); + + reader + .read_to_end(&mut output) + .await + .expect("error after the requested bytes should not trigger a redundant chunk resume"); + assert_eq!(output, b"0123"); + assert!( + transport + .fresh_chunk_requests + .lock() + .expect("fresh chunk request lock should not be poisoned") + .is_empty() + ); + } + + #[tokio::test] + async fn remote_reader_retries_at_most_once_and_preserves_non_retryable_errors() { + let transport = Arc::new(ResumeTransport::with_read_steps(vec![ResumeReadStep::PartialThenReset(b"456".to_vec())])); + let mut reader = RetryingRemoteReader::new_with_timeouts( + resume_step_reader(ResumeReadStep::PartialThenReset(b"0123".to_vec())), + transport.clone(), + resume_request(7), + None, + None, + ); + let error = reader + .read_to_end(&mut Vec::new()) + .await + .expect_err("second reset must not retry"); + assert_eq!(error.kind(), std_io::ErrorKind::ConnectionReset); + assert_eq!( + transport + .read_requests + .lock() + .expect("read request lock should not be poisoned") + .len(), + 1 + ); + + let transport = Arc::new(ResumeTransport::default()); + let reader = PartialThenErrorReader { + cursor: Cursor::new(b"data".to_vec()), + error: Some(io::Error::new(std_io::ErrorKind::PermissionDenied, "permission denied")), + }; + let mut reader = + RetryingRemoteReader::new_with_timeouts(Box::new(reader), transport.clone(), resume_request(4), None, None); + let error = reader + .read_to_end(&mut Vec::new()) + .await + .expect_err("non-retryable errors must not retry"); + assert_eq!(error.kind(), std_io::ErrorKind::PermissionDenied); + assert!( + transport + .read_requests + .lock() + .expect("read request lock should not be poisoned") + .is_empty() + ); + } + + #[test] + fn resumed_read_request_checks_large_offsets() { + let request = ReadStreamRequest { + offset: usize::MAX - 1, + length: 0, + ..resume_request(0) + }; + assert!(resumed_read_request(&request, 2).is_err()); + + let request = resume_request(4); + assert!(resumed_read_request(&request, 5).is_err()); + } + fn init_tracing(filter_level: Level) { INIT.call_once(|| { let _ = tracing_subscriber::fmt() @@ -4509,7 +5670,7 @@ mod tests { assert_eq!(request.path, "object/part.1"); assert_eq!(request.offset, 7); assert_eq!(request.length, 11); - assert_eq!(request.stall_timeout, Some(get_object_disk_read_timeout())); + assert_eq!(request.stall_timeout, remote_read_timeouts(get_object_disk_read_timeout()).body_stall); } other => panic!("expected read transport call, got {other:?}"), } diff --git a/crates/rio/src/http_reader.rs b/crates/rio/src/http_reader.rs index 96fd0971f..93f61a15c 100644 --- a/crates/rio/src/http_reader.rs +++ b/crates/rio/src/http_reader.rs @@ -715,6 +715,13 @@ async fn get_http_client(url: &str) -> io::Result { Ok(cached.client_for(disable_proxy)) } +async fn get_fresh_http_client(url: &str) -> io::Result { + let tuning = internode_http_client_tuning(); + let disable_proxy = should_disable_proxy_for_url(url, tuning); + let outbound_tls = crate::http_runtime_sources::outbound_tls_state().await; + build_http_client(disable_proxy, tuning, &outbound_tls).await +} + fn internode_request_context(method: &Method, url: &str, operation: Option<&'static str>) -> InternodeHttpRequestContext { let target = reqwest::Url::parse(url) .ok() @@ -962,6 +969,28 @@ impl HttpReader { Self::with_capacity_and_stall_timeout(url, method, headers, body, 0, stall_timeout).await } + pub async fn new_fresh_connection_with_stall_timeout( + url: String, + method: Method, + headers: HeaderMap, + body: Option>, + stall_timeout: Option, + ) -> io::Result { + let init = Self::open(&url, &method, &headers, body, stall_timeout, true).await?; + Ok(Self { + inner: StreamReader::new(init.stream), + url, + method, + headers, + track_internode_metrics: init.track_internode_metrics, + internode_operation: init.internode_operation, + stall_timer: None, + stall_timeout: init.stall_timeout, + request_started: init.request_started, + duration_recorded: false, + }) + } + /// Create a new HttpReader from a URL. The request is performed immediately. pub async fn with_capacity( url: String, @@ -981,7 +1010,7 @@ impl HttpReader { _read_buf_size: usize, stall_timeout: Option, ) -> io::Result { - let init = Self::open(&url, &method, &headers, body, stall_timeout).await?; + let init = Self::open(&url, &method, &headers, body, stall_timeout, false).await?; Ok(Self { inner: StreamReader::new(init.stream), url, @@ -1002,10 +1031,16 @@ impl HttpReader { headers: &HeaderMap, body: Option>, stall_timeout: Option, + force_fresh_connection: bool, ) -> io::Result { let track_internode_metrics = is_internode_rpc_url(url); let internode_operation = internode_rpc_operation(url); - let client = get_http_client(url).await.inspect_err(|_| { + let client = if force_fresh_connection { + get_fresh_http_client(url).await + } else { + get_http_client(url).await + } + .inspect_err(|_| { record_internode_error(track_internode_metrics, internode_operation); })?; let mut request: RequestBuilder = client.request(method.clone(), url).headers(headers.clone()); @@ -1121,7 +1156,28 @@ impl HttpChunkReader { body: Option>, stall_timeout: Option, ) -> io::Result { - let init = HttpReader::open(&url, &method, &headers, body, stall_timeout).await?; + let init = HttpReader::open(&url, &method, &headers, body, stall_timeout, false).await?; + Ok(Self { + inner: init.stream, + current: None, + track_internode_metrics: init.track_internode_metrics, + internode_operation: init.internode_operation, + stall_timer: None, + stall_timeout: init.stall_timeout, + request_started: init.request_started, + duration_recorded: false, + consecutive_empty_chunks: 0, + }) + } + + pub async fn new_fresh_connection_with_stall_timeout( + url: String, + method: Method, + headers: HeaderMap, + body: Option>, + stall_timeout: Option, + ) -> io::Result { + let init = HttpReader::open(&url, &method, &headers, body, stall_timeout, true).await?; Ok(Self { inner: init.stream, current: None,