mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-16 09:58:21 +00:00
Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3b3eb895d4 | |||
| 0d2dc23854 | |||
| 9699579ae1 | |||
| ca96da9fa2 | |||
| b1268173fb | |||
| c98d6b58a1 | |||
| b18ccefd1b |
@@ -233,11 +233,17 @@ pub struct NsScannerCapabilityRequest {
|
|||||||
#[async_trait]
|
#[async_trait]
|
||||||
pub trait InternodeDataTransport: Send + Sync + std::fmt::Debug {
|
pub trait InternodeDataTransport: Send + Sync + std::fmt::Debug {
|
||||||
async fn open_read(&self, request: ReadStreamRequest) -> Result<FileReader>;
|
async fn open_read(&self, request: ReadStreamRequest) -> Result<FileReader>;
|
||||||
|
async fn open_read_fresh(&self, request: ReadStreamRequest) -> Result<FileReader> {
|
||||||
|
self.open_read(request).await
|
||||||
|
}
|
||||||
/// Opens an owned-chunk stream when this transport can retain receive-buffer
|
/// Opens an owned-chunk stream when this transport can retain receive-buffer
|
||||||
/// ownership. `None` preserves the established `open_read` fallback.
|
/// ownership. `None` preserves the established `open_read` fallback.
|
||||||
async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result<Option<ChunkReaderBox>> {
|
async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result<Option<ChunkReaderBox>> {
|
||||||
Ok(None)
|
Ok(None)
|
||||||
}
|
}
|
||||||
|
async fn open_read_chunks_fresh(&self, request: ReadStreamRequest) -> Result<Option<ChunkReaderBox>> {
|
||||||
|
self.open_read_chunks(request).await
|
||||||
|
}
|
||||||
async fn open_write(&self, request: WriteStreamRequest) -> Result<FileWriter>;
|
async fn open_write(&self, request: WriteStreamRequest) -> Result<FileWriter>;
|
||||||
async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result<FileReader>;
|
async fn open_walk_dir(&self, request: WalkDirStreamRequest) -> Result<FileReader>;
|
||||||
async fn open_ns_scanner(&self, _request: NsScannerStreamRequest) -> Result<FileReader> {
|
async fn open_ns_scanner(&self, _request: NsScannerStreamRequest) -> Result<FileReader> {
|
||||||
@@ -269,6 +275,15 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport {
|
|||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn open_read_fresh(&self, request: ReadStreamRequest) -> Result<FileReader> {
|
||||||
|
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<Option<ChunkReaderBox>> {
|
async fn open_read_chunks(&self, request: ReadStreamRequest) -> Result<Option<ChunkReaderBox>> {
|
||||||
let url = build_read_file_stream_url(&request);
|
let url = build_read_file_stream_url(&request);
|
||||||
let mut headers = json_headers();
|
let mut headers = json_headers();
|
||||||
@@ -278,6 +293,16 @@ impl InternodeDataTransport for TcpHttpInternodeDataTransport {
|
|||||||
)))
|
)))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn open_read_chunks_fresh(&self, request: ReadStreamRequest) -> Result<Option<ChunkReaderBox>> {
|
||||||
|
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<FileWriter> {
|
async fn open_write(&self, request: WriteStreamRequest) -> Result<FileWriter> {
|
||||||
let server_epoch = self.put_file_auth_capability(&request.endpoint).await?;
|
let server_epoch = self.put_file_auth_capability(&request.endpoint).await?;
|
||||||
let nonce = server_epoch.map(|_| Uuid::new_v4());
|
let nonce = server_epoch.map(|_| Uuid::new_v4());
|
||||||
|
|||||||
@@ -57,15 +57,17 @@ use serde::{Serialize, de::DeserializeOwned};
|
|||||||
use std::{
|
use std::{
|
||||||
io::Cursor,
|
io::Cursor,
|
||||||
path::PathBuf,
|
path::PathBuf,
|
||||||
|
pin::Pin,
|
||||||
sync::{
|
sync::{
|
||||||
Arc,
|
Arc,
|
||||||
atomic::{AtomicU32, Ordering},
|
atomic::{AtomicU32, Ordering},
|
||||||
},
|
},
|
||||||
|
task::{Context, Poll},
|
||||||
time::Duration,
|
time::Duration,
|
||||||
};
|
};
|
||||||
use tokio::time;
|
use tokio::time;
|
||||||
use tokio::{
|
use tokio::{
|
||||||
io::{self, AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt},
|
io::{self, AsyncRead, AsyncReadExt, AsyncWrite, AsyncWriteExt, ReadBuf},
|
||||||
net::TcpStream,
|
net::TcpStream,
|
||||||
time::timeout,
|
time::timeout,
|
||||||
};
|
};
|
||||||
@@ -214,6 +216,231 @@ where
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn is_retryable_remote_body_error(error: &io::Error) -> bool {
|
||||||
|
if error
|
||||||
|
.get_ref()
|
||||||
|
.and_then(|source| source.downcast_ref::<rustfs_rio::BodyStalled>())
|
||||||
|
.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<ReadStreamRequest> {
|
||||||
|
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()
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
type ReadResumeFuture = tokio::task::JoinHandle<Result<FileReader>>;
|
||||||
|
|
||||||
|
struct RetryingRemoteReader {
|
||||||
|
reader: Option<FileReader>,
|
||||||
|
transport: Arc<dyn InternodeDataTransport>,
|
||||||
|
request: ReadStreamRequest,
|
||||||
|
emitted: usize,
|
||||||
|
retried: bool,
|
||||||
|
resume: Option<ReadResumeFuture>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RetryingRemoteReader {
|
||||||
|
fn new(reader: FileReader, transport: Arc<dyn InternodeDataTransport>, request: ReadStreamRequest) -> Self {
|
||||||
|
Self {
|
||||||
|
reader: Some(reader),
|
||||||
|
transport,
|
||||||
|
request,
|
||||||
|
emitted: 0,
|
||||||
|
retried: false,
|
||||||
|
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 transport = Arc::clone(&self.transport);
|
||||||
|
self.resume = Some(tokio::spawn(async move { transport.open_read_fresh(request).await }));
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AsyncRead for RetryingRemoteReader {
|
||||||
|
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<io::Result<()>> {
|
||||||
|
loop {
|
||||||
|
if let Some(resume) = self.resume.as_mut() {
|
||||||
|
match Pin::new(resume).poll(cx) {
|
||||||
|
Poll::Pending => return Poll::Pending,
|
||||||
|
Poll::Ready(Ok(Ok(reader))) => {
|
||||||
|
self.resume = None;
|
||||||
|
self.reader = Some(reader);
|
||||||
|
}
|
||||||
|
Poll::Ready(Ok(Err(error))) => {
|
||||||
|
self.resume = None;
|
||||||
|
return Poll::Ready(Err(io::Error::other(error)));
|
||||||
|
}
|
||||||
|
Poll::Ready(Err(error)) => {
|
||||||
|
self.resume = None;
|
||||||
|
return Poll::Ready(Err(io::Error::other(error)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let Some(reader) = self.reader.as_mut() else {
|
||||||
|
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"))),
|
||||||
|
};
|
||||||
|
return Poll::Ready(Ok(()));
|
||||||
|
}
|
||||||
|
Poll::Ready(Err(error)) if !self.retried && is_retryable_remote_body_error(&error) => {
|
||||||
|
self.retried = true;
|
||||||
|
if let Err(resume_error) = self.start_resume() {
|
||||||
|
return Poll::Ready(Err(resume_error));
|
||||||
|
}
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
Poll::Ready(Err(error)) => return Poll::Ready(Err(error)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type ChunkResumeFuture = tokio::task::JoinHandle<Result<Option<rustfs_rio::ChunkReaderBox>>>;
|
||||||
|
|
||||||
|
struct RetryingRemoteChunkReader {
|
||||||
|
reader: Option<rustfs_rio::ChunkReaderBox>,
|
||||||
|
transport: Arc<dyn InternodeDataTransport>,
|
||||||
|
request: ReadStreamRequest,
|
||||||
|
emitted: usize,
|
||||||
|
retried: bool,
|
||||||
|
resume: Option<ChunkResumeFuture>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RetryingRemoteChunkReader {
|
||||||
|
fn new(reader: rustfs_rio::ChunkReaderBox, transport: Arc<dyn InternodeDataTransport>, request: ReadStreamRequest) -> Self {
|
||||||
|
Self {
|
||||||
|
reader: Some(reader),
|
||||||
|
transport,
|
||||||
|
request,
|
||||||
|
emitted: 0,
|
||||||
|
retried: false,
|
||||||
|
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 transport = Arc::clone(&self.transport);
|
||||||
|
self.resume = Some(tokio::spawn(async move { transport.open_read_chunks_fresh(request).await }));
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AsyncRead for RetryingRemoteChunkReader {
|
||||||
|
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(()));
|
||||||
|
}
|
||||||
|
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<io::Result<Option<Bytes>>> {
|
||||||
|
loop {
|
||||||
|
if let Some(resume) = self.resume.as_mut() {
|
||||||
|
match Pin::new(resume).poll(cx) {
|
||||||
|
Poll::Pending => return Poll::Pending,
|
||||||
|
Poll::Ready(Ok(Ok(Some(reader)))) => {
|
||||||
|
self.resume = None;
|
||||||
|
self.reader = Some(reader);
|
||||||
|
}
|
||||||
|
Poll::Ready(Ok(Ok(None))) => {
|
||||||
|
self.resume = None;
|
||||||
|
self.reader = None;
|
||||||
|
return Poll::Ready(Err(io::Error::other("remote resume transport did not provide a chunk reader")));
|
||||||
|
}
|
||||||
|
Poll::Ready(Ok(Err(error))) => {
|
||||||
|
self.resume = None;
|
||||||
|
return Poll::Ready(Err(io::Error::other(error)));
|
||||||
|
}
|
||||||
|
Poll::Ready(Err(error)) => {
|
||||||
|
self.resume = None;
|
||||||
|
return Poll::Ready(Err(io::Error::other(error)));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let Some(reader) = self.reader.as_mut() else {
|
||||||
|
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"))),
|
||||||
|
};
|
||||||
|
return Poll::Ready(Ok(Some(chunk)));
|
||||||
|
}
|
||||||
|
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;
|
||||||
|
if let Err(resume_error) = self.start_resume() {
|
||||||
|
return Poll::Ready(Err(resume_error));
|
||||||
|
}
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
Poll::Ready(Err(error)) => return Poll::Ready(Err(error)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub struct RemoteDisk {
|
pub struct RemoteDisk {
|
||||||
pub id: Mutex<Option<Uuid>>,
|
pub id: Mutex<Option<Uuid>>,
|
||||||
@@ -2469,7 +2696,7 @@ impl DiskAPI for RemoteDisk {
|
|||||||
}
|
}
|
||||||
let disk = self.disk_ref().await;
|
let disk = self.disk_ref().await;
|
||||||
let stall_timeout = get_object_disk_read_timeout();
|
let stall_timeout = get_object_disk_read_timeout();
|
||||||
self.open_read_with_retry(ReadStreamRequest {
|
let request = ReadStreamRequest {
|
||||||
endpoint: self.endpoint.grid_host(),
|
endpoint: self.endpoint.grid_host(),
|
||||||
disk,
|
disk,
|
||||||
volume: volume.to_string(),
|
volume: volume.to_string(),
|
||||||
@@ -2477,8 +2704,9 @@ impl DiskAPI for RemoteDisk {
|
|||||||
offset,
|
offset,
|
||||||
length,
|
length,
|
||||||
stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout),
|
stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout),
|
||||||
})
|
};
|
||||||
.await
|
let reader = self.open_read_with_retry(request.clone()).await?;
|
||||||
|
Ok(Box::new(RetryingRemoteReader::new(reader, Arc::clone(&self.data_transport), request)))
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn read_file_stream_chunks(
|
async fn read_file_stream_chunks(
|
||||||
@@ -2493,7 +2721,7 @@ impl DiskAPI for RemoteDisk {
|
|||||||
}
|
}
|
||||||
let disk = self.disk_ref().await;
|
let disk = self.disk_ref().await;
|
||||||
let stall_timeout = get_object_disk_read_timeout();
|
let stall_timeout = get_object_disk_read_timeout();
|
||||||
self.open_read_chunks_with_retry(ReadStreamRequest {
|
let request = ReadStreamRequest {
|
||||||
endpoint: self.endpoint.grid_host(),
|
endpoint: self.endpoint.grid_host(),
|
||||||
disk,
|
disk,
|
||||||
volume: volume.to_string(),
|
volume: volume.to_string(),
|
||||||
@@ -2501,8 +2729,12 @@ impl DiskAPI for RemoteDisk {
|
|||||||
offset,
|
offset,
|
||||||
length,
|
length,
|
||||||
stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout),
|
stall_timeout: (!stall_timeout.is_zero()).then_some(stall_timeout),
|
||||||
})
|
};
|
||||||
.await
|
let reader = self.open_read_chunks_with_retry(request.clone()).await?;
|
||||||
|
Ok(reader.map(|reader| {
|
||||||
|
Box::new(RetryingRemoteChunkReader::new(reader, Arc::clone(&self.data_transport), request))
|
||||||
|
as rustfs_rio::ChunkReaderBox
|
||||||
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Buffered read for remote disks.
|
/// Buffered read for remote disks.
|
||||||
@@ -4123,6 +4355,278 @@ mod tests {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone)]
|
||||||
|
enum ResumeReadStep {
|
||||||
|
PartialThenReset(Vec<u8>),
|
||||||
|
Data(Vec<u8>),
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Default)]
|
||||||
|
struct ResumeTransport {
|
||||||
|
read_steps: Mutex<Vec<ResumeReadStep>>,
|
||||||
|
chunk_steps: Mutex<Vec<ResumeReadStep>>,
|
||||||
|
read_requests: Mutex<Vec<ReadStreamRequest>>,
|
||||||
|
chunk_requests: Mutex<Vec<ReadStreamRequest>>,
|
||||||
|
fresh_read_requests: Mutex<Vec<ReadStreamRequest>>,
|
||||||
|
fresh_chunk_requests: Mutex<Vec<ReadStreamRequest>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ResumeTransport {
|
||||||
|
fn with_read_steps(read_steps: Vec<ResumeReadStep>) -> Self {
|
||||||
|
Self {
|
||||||
|
read_steps: Mutex::new(read_steps),
|
||||||
|
..Self::default()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn with_chunk_steps(chunk_steps: Vec<ResumeReadStep>) -> Self {
|
||||||
|
Self {
|
||||||
|
chunk_steps: Mutex::new(chunk_steps),
|
||||||
|
..Self::default()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct ChunkPartialThenErrorReader {
|
||||||
|
data: Option<Bytes>,
|
||||||
|
error: Option<io::Error>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl rustfs_rio::ChunkReader for ChunkPartialThenErrorReader {
|
||||||
|
fn poll_read_chunk(mut self: Pin<&mut Self>, _cx: &mut Context<'_>, max: usize) -> Poll<io::Result<Option<Bytes>>> {
|
||||||
|
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<io::Result<()>> {
|
||||||
|
Poll::Ready(Err(io::Error::other("chunk reader must use chunk handoff")))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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<FileReader> {
|
||||||
|
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<FileReader> {
|
||||||
|
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<Option<rustfs_rio::ChunkReaderBox>> {
|
||||||
|
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<Option<rustfs_rio::ChunkReaderBox>> {
|
||||||
|
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<FileWriter> {
|
||||||
|
panic!("open_write should not be used in remote read resume tests");
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn open_walk_dir(&self, _request: WalkDirStreamRequest) -> Result<FileReader> {
|
||||||
|
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(reader, transport.clone(), request);
|
||||||
|
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(reader, transport.clone(), request);
|
||||||
|
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
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[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(
|
||||||
|
resume_step_reader(ResumeReadStep::PartialThenReset(b"0123".to_vec())),
|
||||||
|
transport.clone(),
|
||||||
|
resume_request(7),
|
||||||
|
);
|
||||||
|
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(Box::new(reader), transport.clone(), resume_request(4));
|
||||||
|
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) {
|
fn init_tracing(filter_level: Level) {
|
||||||
INIT.call_once(|| {
|
INIT.call_once(|| {
|
||||||
let _ = tracing_subscriber::fmt()
|
let _ = tracing_subscriber::fmt()
|
||||||
|
|||||||
@@ -1093,6 +1093,8 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
|||||||
let (rd, wd) = tokio::io::duplex(duplex_buffer_size);
|
let (rd, wd) = tokio::io::duplex(duplex_buffer_size);
|
||||||
debug!(bucket, object, duplex_buffer_size, "Created duplex pipe for object data transfer");
|
debug!(bucket, object, duplex_buffer_size, "Created duplex pipe for object data transfer");
|
||||||
|
|
||||||
|
let (producer_terminal_tx, producer_terminal_rx) = tokio::sync::oneshot::channel();
|
||||||
|
let rd = LegacyDuplexProducerReader::new(rd, producer_terminal_rx);
|
||||||
let (mut reader, offset, length) =
|
let (mut reader, offset, length) =
|
||||||
get_object_reader_with_context(&self.ctx, Box::new(rd), range, &object_info, opts, &h).await?;
|
get_object_reader_with_context(&self.ctx, Box::new(rd), range, &object_info, opts, &h).await?;
|
||||||
// Carry the hook probe result so the app layer skips its now-redundant
|
// Carry the hook probe result so the app layer skips its now-redundant
|
||||||
@@ -1114,7 +1116,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
|||||||
// `get_object_with_fileinfo` also waits on `writer`, so an outer timeout
|
// `get_object_with_fileinfo` also waits on `writer`, so an outer timeout
|
||||||
// would incorrectly treat downstream backpressure as disk-read latency.
|
// would incorrectly treat downstream backpressure as disk-read latency.
|
||||||
// Disk read timeouts must be enforced at the actual disk I/O operations.
|
// Disk read timeouts must be enforced at the actual disk I/O operations.
|
||||||
if let Err(e) = Self::get_object_with_fileinfo(
|
let producer_result = Self::get_object_with_fileinfo(
|
||||||
&bucket,
|
&bucket,
|
||||||
&object,
|
&object,
|
||||||
erasure_cache,
|
erasure_cache,
|
||||||
@@ -1132,9 +1134,9 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
|||||||
object_class.as_str(),
|
object_class.as_str(),
|
||||||
size_bucket,
|
size_bucket,
|
||||||
)
|
)
|
||||||
.await
|
.await;
|
||||||
{
|
if let Err(e) = &producer_result {
|
||||||
let reason = classify_storage_error(&e);
|
let reason = classify_storage_error(e);
|
||||||
if reason == GetObjectFailureReason::DownstreamClosed {
|
if reason == GetObjectFailureReason::DownstreamClosed {
|
||||||
debug!(
|
debug!(
|
||||||
event = EVENT_SET_DISK_WRITE,
|
event = EVENT_SET_DISK_WRITE,
|
||||||
@@ -1173,6 +1175,7 @@ impl crate::storage_api_contracts::object::ObjectIO for SetDisks {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
let _ = producer_terminal_tx.send(producer_result.map(|_| ()));
|
||||||
});
|
});
|
||||||
|
|
||||||
Ok(reader)
|
Ok(reader)
|
||||||
@@ -2318,6 +2321,159 @@ impl<R: AsyncRead + Unpin> AsyncRead for TransitionUploadReader<R> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
struct LegacyDuplexProducerReader<R> {
|
||||||
|
inner: R,
|
||||||
|
terminal: Option<tokio::sync::oneshot::Receiver<Result<()>>>,
|
||||||
|
inner_eof: bool,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<R> LegacyDuplexProducerReader<R> {
|
||||||
|
fn new(inner: R, terminal: tokio::sync::oneshot::Receiver<Result<()>>) -> Self {
|
||||||
|
Self {
|
||||||
|
inner,
|
||||||
|
terminal: Some(terminal),
|
||||||
|
inner_eof: false,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<R: AsyncRead + Unpin> AsyncRead for LegacyDuplexProducerReader<R> {
|
||||||
|
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||||
|
if !self.inner_eof {
|
||||||
|
let before = buf.filled().len();
|
||||||
|
match Pin::new(&mut self.inner).poll_read(cx, buf) {
|
||||||
|
Poll::Pending => return Poll::Pending,
|
||||||
|
Poll::Ready(Err(err)) => return Poll::Ready(Err(err)),
|
||||||
|
Poll::Ready(Ok(())) if buf.filled().len() > before => return Poll::Ready(Ok(())),
|
||||||
|
Poll::Ready(Ok(())) => {
|
||||||
|
self.inner_eof = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
let Some(terminal) = self.terminal.as_mut() else {
|
||||||
|
return Poll::Ready(Ok(()));
|
||||||
|
};
|
||||||
|
match Pin::new(terminal).poll(cx) {
|
||||||
|
Poll::Pending => Poll::Pending,
|
||||||
|
Poll::Ready(Ok(Ok(()))) => {
|
||||||
|
self.terminal = None;
|
||||||
|
Poll::Ready(Ok(()))
|
||||||
|
}
|
||||||
|
Poll::Ready(Ok(Err(err))) => {
|
||||||
|
self.terminal = None;
|
||||||
|
Poll::Ready(Err(std::io::Error::other(err)))
|
||||||
|
}
|
||||||
|
Poll::Ready(Err(_)) => {
|
||||||
|
self.terminal = None;
|
||||||
|
Poll::Ready(Err(std::io::Error::other(StorageError::Unexpected)))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod legacy_duplex_producer_reader_tests {
|
||||||
|
use super::*;
|
||||||
|
use tokio::io::{AsyncReadExt, AsyncWriteExt};
|
||||||
|
|
||||||
|
fn storage_error_source(error: &std::io::Error) -> &StorageError {
|
||||||
|
error
|
||||||
|
.get_ref()
|
||||||
|
.and_then(|source| source.downcast_ref::<StorageError>())
|
||||||
|
.expect("legacy duplex terminal error should retain StorageError source")
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn legacy_duplex_reader_allows_clean_completion() {
|
||||||
|
let (mut writer, reader) = tokio::io::duplex(64);
|
||||||
|
let (terminal_tx, terminal_rx) = tokio::sync::oneshot::channel();
|
||||||
|
writer
|
||||||
|
.write_all(b"complete")
|
||||||
|
.await
|
||||||
|
.expect("duplex write should fit in buffer");
|
||||||
|
drop(writer);
|
||||||
|
terminal_tx.send(Ok(())).expect("terminal receiver should remain installed");
|
||||||
|
|
||||||
|
let mut reader = LegacyDuplexProducerReader::new(reader, terminal_rx);
|
||||||
|
let mut out = Vec::new();
|
||||||
|
reader
|
||||||
|
.read_to_end(&mut out)
|
||||||
|
.await
|
||||||
|
.expect("clean producer completion should surface clean EOF");
|
||||||
|
|
||||||
|
assert_eq!(out, b"complete");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn legacy_duplex_reader_surfaces_terminal_error_after_partial_data() {
|
||||||
|
let (mut writer, reader) = tokio::io::duplex(64);
|
||||||
|
let (terminal_tx, terminal_rx) = tokio::sync::oneshot::channel();
|
||||||
|
writer.write_all(b"partial").await.expect("duplex write should fit in buffer");
|
||||||
|
drop(writer);
|
||||||
|
terminal_tx
|
||||||
|
.send(Err(StorageError::FileCorrupt))
|
||||||
|
.expect("terminal receiver should remain installed");
|
||||||
|
|
||||||
|
let mut reader = LegacyDuplexProducerReader::new(reader, terminal_rx);
|
||||||
|
let mut out = Vec::new();
|
||||||
|
let err = reader
|
||||||
|
.read_to_end(&mut out)
|
||||||
|
.await
|
||||||
|
.expect_err("terminal producer error must not become clean EOF");
|
||||||
|
|
||||||
|
assert_eq!(out, b"partial");
|
||||||
|
assert!(matches!(storage_error_source(&err), StorageError::FileCorrupt));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn legacy_duplex_reader_surfaces_terminal_error_after_declared_length() {
|
||||||
|
let (mut writer, reader) = tokio::io::duplex(64);
|
||||||
|
let (terminal_tx, terminal_rx) = tokio::sync::oneshot::channel();
|
||||||
|
writer.write_all(b"exact").await.expect("duplex write should fit in buffer");
|
||||||
|
drop(writer);
|
||||||
|
terminal_tx
|
||||||
|
.send(Err(StorageError::Io(std::io::Error::new(
|
||||||
|
std::io::ErrorKind::ConnectionReset,
|
||||||
|
"remote body reset after final byte",
|
||||||
|
))))
|
||||||
|
.expect("terminal receiver should remain installed");
|
||||||
|
|
||||||
|
let reader = LegacyDuplexProducerReader::new(reader, terminal_rx);
|
||||||
|
let mut reader =
|
||||||
|
HashReader::from_stream(reader, 5, 5, None, None, false).expect("hash reader should accept exact declared length");
|
||||||
|
let mut out = Vec::new();
|
||||||
|
let err = reader
|
||||||
|
.read_to_end(&mut out)
|
||||||
|
.await
|
||||||
|
.expect_err("producer terminal error after the declared length must still fail");
|
||||||
|
|
||||||
|
assert_eq!(out, b"exact");
|
||||||
|
assert!(
|
||||||
|
matches!(storage_error_source(&err), StorageError::Io(io_error) if io_error.kind() == std::io::ErrorKind::ConnectionReset)
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn legacy_duplex_reader_fails_closed_when_terminal_channel_closes() {
|
||||||
|
let (mut writer, reader) = tokio::io::duplex(64);
|
||||||
|
let (terminal_tx, terminal_rx) = tokio::sync::oneshot::channel::<Result<()>>();
|
||||||
|
writer.write_all(b"body").await.expect("duplex write should fit in buffer");
|
||||||
|
drop(writer);
|
||||||
|
drop(terminal_tx);
|
||||||
|
|
||||||
|
let mut reader = LegacyDuplexProducerReader::new(reader, terminal_rx);
|
||||||
|
let mut out = Vec::new();
|
||||||
|
let err = reader
|
||||||
|
.read_to_end(&mut out)
|
||||||
|
.await
|
||||||
|
.expect_err("producer disappearance must fail closed");
|
||||||
|
|
||||||
|
assert_eq!(out, b"body");
|
||||||
|
assert!(matches!(storage_error_source(&err), StorageError::Unexpected));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
struct TransitionUploadWriter<W> {
|
struct TransitionUploadWriter<W> {
|
||||||
inner: W,
|
inner: W,
|
||||||
produced: u64,
|
produced: u64,
|
||||||
|
|||||||
+111
-11
@@ -138,6 +138,12 @@ impl std::fmt::Display for InternodeHttpErrorKind {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[derive(thiserror::Error, Debug, Clone, Copy, Eq, PartialEq)]
|
||||||
|
#[error("internode body stalled for {timeout:?}")]
|
||||||
|
pub struct BodyStalled {
|
||||||
|
pub timeout: Duration,
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone, Eq, PartialEq)]
|
#[derive(Debug, Clone, Eq, PartialEq)]
|
||||||
pub struct InternodeHttpRequestContext {
|
pub struct InternodeHttpRequestContext {
|
||||||
method: String,
|
method: String,
|
||||||
@@ -271,6 +277,10 @@ pub fn internode_http_timeout_error(method: &Method, url: &str) -> io::Error {
|
|||||||
internode_kind_error(method, url, internode_rpc_operation(url), InternodeHttpErrorKind::ConnectTimeout)
|
internode_kind_error(method, url, internode_rpc_operation(url), InternodeHttpErrorKind::ConnectTimeout)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn body_stalled_error(stall_timeout: Duration) -> io::Error {
|
||||||
|
Error::new(io::ErrorKind::TimedOut, BodyStalled { timeout: stall_timeout })
|
||||||
|
}
|
||||||
|
|
||||||
/// Clone an internode HTTP I/O error while retaining its structured classification.
|
/// Clone an internode HTTP I/O error while retaining its structured classification.
|
||||||
///
|
///
|
||||||
/// The underlying transport source is intentionally omitted because it is not
|
/// The underlying transport source is intentionally omitted because it is not
|
||||||
@@ -705,6 +715,13 @@ async fn get_http_client(url: &str) -> io::Result<Client> {
|
|||||||
Ok(cached.client_for(disable_proxy))
|
Ok(cached.client_for(disable_proxy))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn get_fresh_http_client(url: &str) -> io::Result<Client> {
|
||||||
|
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 {
|
fn internode_request_context(method: &Method, url: &str, operation: Option<&'static str>) -> InternodeHttpRequestContext {
|
||||||
let target = reqwest::Url::parse(url)
|
let target = reqwest::Url::parse(url)
|
||||||
.ok()
|
.ok()
|
||||||
@@ -952,6 +969,28 @@ impl HttpReader {
|
|||||||
Self::with_capacity_and_stall_timeout(url, method, headers, body, 0, stall_timeout).await
|
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<Vec<u8>>,
|
||||||
|
stall_timeout: Option<Duration>,
|
||||||
|
) -> io::Result<Self> {
|
||||||
|
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.
|
/// Create a new HttpReader from a URL. The request is performed immediately.
|
||||||
pub async fn with_capacity(
|
pub async fn with_capacity(
|
||||||
url: String,
|
url: String,
|
||||||
@@ -971,7 +1010,7 @@ impl HttpReader {
|
|||||||
_read_buf_size: usize,
|
_read_buf_size: usize,
|
||||||
stall_timeout: Option<Duration>,
|
stall_timeout: Option<Duration>,
|
||||||
) -> io::Result<Self> {
|
) -> io::Result<Self> {
|
||||||
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 {
|
Ok(Self {
|
||||||
inner: StreamReader::new(init.stream),
|
inner: StreamReader::new(init.stream),
|
||||||
url,
|
url,
|
||||||
@@ -992,10 +1031,16 @@ impl HttpReader {
|
|||||||
headers: &HeaderMap,
|
headers: &HeaderMap,
|
||||||
body: Option<Vec<u8>>,
|
body: Option<Vec<u8>>,
|
||||||
stall_timeout: Option<Duration>,
|
stall_timeout: Option<Duration>,
|
||||||
|
force_fresh_connection: bool,
|
||||||
) -> io::Result<HttpReaderInit> {
|
) -> io::Result<HttpReaderInit> {
|
||||||
let track_internode_metrics = is_internode_rpc_url(url);
|
let track_internode_metrics = is_internode_rpc_url(url);
|
||||||
let internode_operation = internode_rpc_operation(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);
|
record_internode_error(track_internode_metrics, internode_operation);
|
||||||
})?;
|
})?;
|
||||||
let mut request: RequestBuilder = client.request(method.clone(), url).headers(headers.clone());
|
let mut request: RequestBuilder = client.request(method.clone(), url).headers(headers.clone());
|
||||||
@@ -1085,10 +1130,7 @@ impl AsyncRead for HttpReader {
|
|||||||
);
|
);
|
||||||
record_internode_stall_timeout(*this.track_internode_metrics, *this.internode_operation);
|
record_internode_stall_timeout(*this.track_internode_metrics, *this.internode_operation);
|
||||||
record_internode_error(*this.track_internode_metrics, *this.internode_operation);
|
record_internode_error(*this.track_internode_metrics, *this.internode_operation);
|
||||||
Poll::Ready(Err(Error::new(
|
Poll::Ready(Err(body_stalled_error(stall_timeout)))
|
||||||
io::ErrorKind::TimedOut,
|
|
||||||
"HttpReader stall timeout: no data received before deadline",
|
|
||||||
)))
|
|
||||||
} else {
|
} else {
|
||||||
Poll::Pending
|
Poll::Pending
|
||||||
}
|
}
|
||||||
@@ -1114,7 +1156,28 @@ impl HttpChunkReader {
|
|||||||
body: Option<Vec<u8>>,
|
body: Option<Vec<u8>>,
|
||||||
stall_timeout: Option<Duration>,
|
stall_timeout: Option<Duration>,
|
||||||
) -> io::Result<Self> {
|
) -> io::Result<Self> {
|
||||||
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<Vec<u8>>,
|
||||||
|
stall_timeout: Option<Duration>,
|
||||||
|
) -> io::Result<Self> {
|
||||||
|
let init = HttpReader::open(&url, &method, &headers, body, stall_timeout, true).await?;
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
inner: init.stream,
|
inner: init.stream,
|
||||||
current: None,
|
current: None,
|
||||||
@@ -1217,10 +1280,7 @@ impl ChunkReader for HttpChunkReader {
|
|||||||
);
|
);
|
||||||
record_internode_stall_timeout(*this.track_internode_metrics, *this.internode_operation);
|
record_internode_stall_timeout(*this.track_internode_metrics, *this.internode_operation);
|
||||||
record_internode_error(*this.track_internode_metrics, *this.internode_operation);
|
record_internode_error(*this.track_internode_metrics, *this.internode_operation);
|
||||||
return Poll::Ready(Err(Error::new(
|
return Poll::Ready(Err(body_stalled_error(stall_timeout)));
|
||||||
io::ErrorKind::TimedOut,
|
|
||||||
"HttpReader stall timeout: no data received before deadline",
|
|
||||||
)));
|
|
||||||
}
|
}
|
||||||
return Poll::Pending;
|
return Poll::Pending;
|
||||||
}
|
}
|
||||||
@@ -2379,6 +2439,46 @@ mod tests {
|
|||||||
Err(err) => err,
|
Err(err) => err,
|
||||||
};
|
};
|
||||||
assert_eq!(err.kind(), io::ErrorKind::TimedOut);
|
assert_eq!(err.kind(), io::ErrorKind::TimedOut);
|
||||||
|
let stalled = err
|
||||||
|
.get_ref()
|
||||||
|
.and_then(|source| source.downcast_ref::<BodyStalled>())
|
||||||
|
.expect("stall timeout should retain typed body-stalled source");
|
||||||
|
assert_eq!(stalled.timeout, Duration::from_millis(20));
|
||||||
|
|
||||||
|
handle.abort();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn http_chunk_reader_stall_timeout_retains_typed_source() {
|
||||||
|
let state = TestState::default();
|
||||||
|
let Some((base_url, handle)) = start_test_server(state).await else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let url = base_url.replace("/stream", "/stall");
|
||||||
|
let mut reader =
|
||||||
|
HttpChunkReader::new_with_stall_timeout(url, Method::GET, HeaderMap::new(), None, Some(Duration::from_millis(20)))
|
||||||
|
.await
|
||||||
|
.expect("chunk reader should open");
|
||||||
|
|
||||||
|
let first = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 64))
|
||||||
|
.await
|
||||||
|
.expect("initial body chunk should arrive")
|
||||||
|
.expect("initial body chunk should not be EOF");
|
||||||
|
assert_eq!(first, b"hello"[..]);
|
||||||
|
|
||||||
|
let err = tokio::time::timeout(
|
||||||
|
Duration::from_secs(1),
|
||||||
|
std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 64)),
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.expect("stall timeout should wake chunk reader")
|
||||||
|
.expect_err("chunk reader should return a timeout error");
|
||||||
|
assert_eq!(err.kind(), io::ErrorKind::TimedOut);
|
||||||
|
let stalled = err
|
||||||
|
.get_ref()
|
||||||
|
.and_then(|source| source.downcast_ref::<BodyStalled>())
|
||||||
|
.expect("chunk stall timeout should retain typed body-stalled source");
|
||||||
|
assert_eq!(stalled.timeout, Duration::from_millis(20));
|
||||||
|
|
||||||
handle.abort();
|
handle.abort();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2173,10 +2173,11 @@ fn get_object_resume_control(ctx: GetObjectResumeContext) -> GetObjectResumeCont
|
|||||||
/// disks" failures keep the existing fail-loud behavior.
|
/// disks" failures keep the existing fail-loud behavior.
|
||||||
fn is_object_relocation_error(err: &std::io::Error) -> bool {
|
fn is_object_relocation_error(err: &std::io::Error) -> bool {
|
||||||
let Some(inner) = err.get_ref() else { return false };
|
let Some(inner) = err.get_ref() else { return false };
|
||||||
matches!(
|
match inner.downcast_ref::<StorageError>() {
|
||||||
inner.downcast_ref::<StorageError>(),
|
Some(StorageError::FileNotFound | StorageError::ObjectNotFound(..) | StorageError::InsufficientReadQuorum(..)) => true,
|
||||||
Some(StorageError::FileNotFound | StorageError::ObjectNotFound(..) | StorageError::InsufficientReadQuorum(..))
|
Some(StorageError::Io(source)) => source.kind() == std::io::ErrorKind::NotFound,
|
||||||
)
|
_ => false,
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Resolve the S3 request-body inter-chunk read timeout from the environment.
|
/// Resolve the S3 request-body inter-chunk read timeout from the environment.
|
||||||
@@ -13006,6 +13007,7 @@ mod tests {
|
|||||||
StorageError::FileNotFound,
|
StorageError::FileNotFound,
|
||||||
StorageError::ObjectNotFound("test-bucket".to_string(), "relocated-object".to_string()),
|
StorageError::ObjectNotFound("test-bucket".to_string(), "relocated-object".to_string()),
|
||||||
StorageError::InsufficientReadQuorum("test-bucket".to_string(), "relocated-object".to_string()),
|
StorageError::InsufficientReadQuorum("test-bucket".to_string(), "relocated-object".to_string()),
|
||||||
|
StorageError::Io(std::io::Error::new(std::io::ErrorKind::NotFound, "relocated shard disappeared")),
|
||||||
] {
|
] {
|
||||||
let reopen_count = Arc::new(AtomicUsize::new(0));
|
let reopen_count = Arc::new(AtomicUsize::new(0));
|
||||||
let control = counting_resume_control(Arc::clone(&reopen_count), |emitted| {
|
let control = counting_resume_control(Arc::clone(&reopen_count), |emitted| {
|
||||||
|
|||||||
@@ -669,6 +669,23 @@ mod tests {
|
|||||||
assert!(api_error.source.is_some());
|
assert!(api_error.source.is_some());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_api_error_from_storage_io_copy_object_terminal_error_stays_internal() {
|
||||||
|
let io_error = IoError::other(StorageError::FileCorrupt);
|
||||||
|
let storage_error: StorageError = io_error.into();
|
||||||
|
assert!(matches!(storage_error, StorageError::FileCorrupt));
|
||||||
|
|
||||||
|
let api_error: ApiError = storage_error.into();
|
||||||
|
|
||||||
|
assert_eq!(api_error.code, S3ErrorCode::InternalError);
|
||||||
|
let source = api_error
|
||||||
|
.source
|
||||||
|
.as_deref()
|
||||||
|
.and_then(|source| source.downcast_ref::<StorageError>())
|
||||||
|
.expect("API error should retain the storage error source");
|
||||||
|
assert!(matches!(source, StorageError::FileCorrupt));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn test_api_error_from_iam_error() {
|
fn test_api_error_from_iam_error() {
|
||||||
let iam_error = rustfs_iam::error::Error::other("IAM test error");
|
let iam_error = rustfs_iam::error::Error::other("IAM test error");
|
||||||
|
|||||||
Reference in New Issue
Block a user