mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-18 18:46:17 +00:00
fix(ecstore): resume remote shard reads once (#6091)
* fix(ecstore): preserve CopyObject producer errors * fix(ecstore): resume remote shard reads once * fix(app): resume preserved relocation I/O errors * fix(ecstore): reserve remote read recovery budget
This commit is contained in:
@@ -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<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
|
||||
/// ownership. `None` preserves the established `open_read` fallback.
|
||||
async fn open_read_chunks(&self, _request: ReadStreamRequest) -> Result<Option<ChunkReaderBox>> {
|
||||
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_walk_dir(&self, request: WalkDirStreamRequest) -> 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>> {
|
||||
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<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> {
|
||||
let server_epoch = self.put_file_auth_capability(&request.endpoint).await?;
|
||||
let nonce = server_epoch.map(|_| Uuid::new_v4());
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -715,6 +715,13 @@ async fn get_http_client(url: &str) -> io::Result<Client> {
|
||||
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 {
|
||||
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<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.
|
||||
pub async fn with_capacity(
|
||||
url: String,
|
||||
@@ -981,7 +1010,7 @@ impl HttpReader {
|
||||
_read_buf_size: usize,
|
||||
stall_timeout: Option<Duration>,
|
||||
) -> 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 {
|
||||
inner: StreamReader::new(init.stream),
|
||||
url,
|
||||
@@ -1002,10 +1031,16 @@ impl HttpReader {
|
||||
headers: &HeaderMap,
|
||||
body: Option<Vec<u8>>,
|
||||
stall_timeout: Option<Duration>,
|
||||
force_fresh_connection: bool,
|
||||
) -> io::Result<HttpReaderInit> {
|
||||
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<Vec<u8>>,
|
||||
stall_timeout: Option<Duration>,
|
||||
) -> 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 {
|
||||
inner: init.stream,
|
||||
current: None,
|
||||
|
||||
Reference in New Issue
Block a user