mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-22 12:26:37 +00:00
perf(ecstore): retain remote shard HTTP chunks (#5991)
* perf(ecstore): retain remote shard HTTP chunks * fix(ecstore): bound remote shard chunk retention * fix(rio): persist empty chunk limit across polls
This commit is contained in:
+358
-16
@@ -52,6 +52,8 @@ const HTTP_VERSION_10_LABEL: &str = "http/1.0";
|
||||
const HTTP_VERSION_11_LABEL: &str = "http/1.1";
|
||||
const HTTP_VERSION_2_LABEL: &str = "h2";
|
||||
const HTTP_VERSION_UNKNOWN_LABEL: &str = "unknown";
|
||||
const MAX_CONSECUTIVE_EMPTY_CHUNKS: usize = 64;
|
||||
const EXCESSIVE_EMPTY_CHUNKS_ERROR: &str = "HTTP body returned too many empty chunks";
|
||||
pub const INTERNODE_DISK_ERROR_HEADER: &str = "x-rustfs-disk-error";
|
||||
pub const INTERNODE_FILE_NOT_FOUND: &str = "file-not-found";
|
||||
pub const INTERNODE_VOLUME_NOT_FOUND: &str = "volume-not-found";
|
||||
@@ -883,6 +885,26 @@ fn internode_status_error(method: &Method, url: &str, operation: Option<&'static
|
||||
InternodeHttpError::new(classified, context).into_io_error()
|
||||
}
|
||||
|
||||
type HttpByteStream = Pin<Box<dyn Stream<Item = std::io::Result<Bytes>> + Send + Sync>>;
|
||||
|
||||
/// An async reader that can also transfer received HTTP body chunks without
|
||||
/// copying their contents into an intermediate caller buffer.
|
||||
pub trait ChunkReader: AsyncRead + Send + Sync + Unpin {
|
||||
/// Returns the next non-empty owned chunk, limited to `max` bytes.
|
||||
/// `None` is EOF.
|
||||
fn poll_read_chunk(self: Pin<&mut Self>, cx: &mut Context<'_>, max: usize) -> Poll<io::Result<Option<Bytes>>>;
|
||||
}
|
||||
|
||||
pub type ChunkReaderBox = Box<dyn ChunkReader>;
|
||||
|
||||
struct HttpReaderInit {
|
||||
stream: HttpByteStream,
|
||||
track_internode_metrics: bool,
|
||||
internode_operation: Option<&'static str>,
|
||||
stall_timeout: Option<Duration>,
|
||||
request_started: Instant,
|
||||
}
|
||||
|
||||
pin_project! {
|
||||
pub struct HttpReader {
|
||||
url:String,
|
||||
@@ -895,7 +917,22 @@ pin_project! {
|
||||
request_started: Instant,
|
||||
duration_recorded: bool,
|
||||
#[pin]
|
||||
inner: StreamReader<Pin<Box<dyn Stream<Item=std::io::Result<Bytes>>+Send+Sync>>, Bytes>,
|
||||
inner: StreamReader<HttpByteStream, Bytes>,
|
||||
}
|
||||
}
|
||||
|
||||
pin_project! {
|
||||
pub struct HttpChunkReader {
|
||||
track_internode_metrics: bool,
|
||||
internode_operation: Option<&'static str>,
|
||||
stall_timeout: Option<Duration>,
|
||||
stall_timer: Option<Pin<Box<Sleep>>>,
|
||||
request_started: Instant,
|
||||
duration_recorded: bool,
|
||||
consecutive_empty_chunks: usize,
|
||||
#[pin]
|
||||
inner: HttpByteStream,
|
||||
current: Option<Bytes>,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -934,12 +971,34 @@ impl HttpReader {
|
||||
_read_buf_size: usize,
|
||||
stall_timeout: Option<Duration>,
|
||||
) -> io::Result<Self> {
|
||||
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 init = Self::open(&url, &method, &headers, body, stall_timeout).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,
|
||||
})
|
||||
}
|
||||
|
||||
async fn open(
|
||||
url: &str,
|
||||
method: &Method,
|
||||
headers: &HeaderMap,
|
||||
body: Option<Vec<u8>>,
|
||||
stall_timeout: Option<Duration>,
|
||||
) -> 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(|_| {
|
||||
record_internode_error(track_internode_metrics, internode_operation);
|
||||
})?;
|
||||
let mut request: RequestBuilder = client.request(method.clone(), url.clone()).headers(headers.clone());
|
||||
let mut request: RequestBuilder = client.request(method.clone(), url).headers(headers.clone());
|
||||
if let Some(body) = body {
|
||||
request = request.body(body);
|
||||
}
|
||||
@@ -949,7 +1008,7 @@ impl HttpReader {
|
||||
record_internode_operation_duration(track_internode_metrics, internode_operation, request_started.elapsed());
|
||||
record_internode_error(track_internode_metrics, internode_operation);
|
||||
record_internode_classified_error(track_internode_metrics, internode_operation, classify_reqwest_error(&e));
|
||||
internode_reqwest_error(&method, &url, internode_operation, e)
|
||||
internode_reqwest_error(method, url, internode_operation, e)
|
||||
})?;
|
||||
|
||||
record_internode_http_version(track_internode_metrics, internode_operation, http_version_metric_label(resp.version()));
|
||||
@@ -959,12 +1018,12 @@ impl HttpReader {
|
||||
record_internode_operation_duration(track_internode_metrics, internode_operation, request_started.elapsed());
|
||||
record_internode_error(track_internode_metrics, internode_operation);
|
||||
record_internode_classified_error(track_internode_metrics, internode_operation, classified.kind);
|
||||
return Err(internode_classified_error(&method, &url, internode_operation, classified));
|
||||
return Err(internode_classified_error(method, url, internode_operation, classified));
|
||||
}
|
||||
|
||||
record_internode_outgoing_request(track_internode_metrics, internode_operation);
|
||||
|
||||
let stream_error_url = url.clone();
|
||||
let stream_error_url = url.to_owned();
|
||||
let stream_error_method = method.clone();
|
||||
let stream = resp.bytes_stream().map_err(move |e| {
|
||||
record_internode_error(track_internode_metrics, internode_operation);
|
||||
@@ -973,17 +1032,12 @@ impl HttpReader {
|
||||
internode_reqwest_body_error(&stream_error_method, &stream_error_url, internode_operation, e)
|
||||
});
|
||||
|
||||
Ok(Self {
|
||||
inner: StreamReader::new(Box::pin(stream)),
|
||||
url,
|
||||
method,
|
||||
headers,
|
||||
Ok(HttpReaderInit {
|
||||
stream: Box::pin(stream),
|
||||
track_internode_metrics,
|
||||
internode_operation,
|
||||
stall_timer: None,
|
||||
stall_timeout,
|
||||
request_started,
|
||||
duration_recorded: false,
|
||||
})
|
||||
}
|
||||
pub fn url(&self) -> &str {
|
||||
@@ -1000,7 +1054,6 @@ impl HttpReader {
|
||||
impl AsyncRead for HttpReader {
|
||||
fn poll_read(self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||
let mut this = self.project();
|
||||
|
||||
let filled_before = buf.filled().len();
|
||||
match this.inner.as_mut().poll_read(cx, buf) {
|
||||
Poll::Ready(Ok(())) => {
|
||||
@@ -1053,6 +1106,129 @@ impl AsyncRead for HttpReader {
|
||||
}
|
||||
}
|
||||
|
||||
impl HttpChunkReader {
|
||||
pub async fn new_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).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,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
fn excessive_empty_chunks_error() -> Error {
|
||||
Error::new(io::ErrorKind::InvalidData, EXCESSIVE_EMPTY_CHUNKS_ERROR)
|
||||
}
|
||||
|
||||
impl AsyncRead for HttpChunkReader {
|
||||
fn poll_read(mut self: Pin<&mut Self>, cx: &mut Context<'_>, buf: &mut ReadBuf<'_>) -> Poll<std::io::Result<()>> {
|
||||
if buf.remaining() == 0 {
|
||||
return Poll::Ready(Ok(()));
|
||||
}
|
||||
match 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(err)) => Poll::Ready(Err(err)),
|
||||
Poll::Pending => Poll::Pending,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl ChunkReader for HttpChunkReader {
|
||||
fn poll_read_chunk(self: Pin<&mut Self>, cx: &mut Context<'_>, max: usize) -> Poll<io::Result<Option<Bytes>>> {
|
||||
if max == 0 {
|
||||
return Poll::Ready(Err(Error::new(io::ErrorKind::InvalidInput, "chunk read limit must be non-zero")));
|
||||
}
|
||||
|
||||
let mut this = self.project();
|
||||
if *this.consecutive_empty_chunks >= MAX_CONSECUTIVE_EMPTY_CHUNKS {
|
||||
return Poll::Ready(Err(excessive_empty_chunks_error()));
|
||||
}
|
||||
loop {
|
||||
if let Some(mut current) = this.current.take() {
|
||||
let take = current.len().min(max);
|
||||
let chunk = current.split_to(take);
|
||||
if !current.is_empty() {
|
||||
*this.current = Some(current);
|
||||
}
|
||||
record_internode_recv_bytes(*this.track_internode_metrics, *this.internode_operation, take);
|
||||
*this.stall_timer = None;
|
||||
return Poll::Ready(Ok(Some(chunk)));
|
||||
}
|
||||
|
||||
match this.inner.as_mut().poll_next(cx) {
|
||||
Poll::Ready(Some(Ok(bytes))) if bytes.is_empty() => {
|
||||
*this.consecutive_empty_chunks += 1;
|
||||
if *this.consecutive_empty_chunks == MAX_CONSECUTIVE_EMPTY_CHUNKS {
|
||||
record_internode_error(*this.track_internode_metrics, *this.internode_operation);
|
||||
return Poll::Ready(Err(excessive_empty_chunks_error()));
|
||||
}
|
||||
}
|
||||
Poll::Ready(Some(Ok(bytes))) => {
|
||||
*this.consecutive_empty_chunks = 0;
|
||||
*this.current = Some(bytes);
|
||||
}
|
||||
Poll::Ready(Some(Err(err))) => {
|
||||
record_internode_operation_duration_once(
|
||||
*this.track_internode_metrics,
|
||||
*this.internode_operation,
|
||||
*this.request_started,
|
||||
this.duration_recorded,
|
||||
);
|
||||
return Poll::Ready(Err(err));
|
||||
}
|
||||
Poll::Ready(None) => {
|
||||
record_internode_operation_duration_once(
|
||||
*this.track_internode_metrics,
|
||||
*this.internode_operation,
|
||||
*this.request_started,
|
||||
this.duration_recorded,
|
||||
);
|
||||
*this.stall_timer = None;
|
||||
return Poll::Ready(Ok(None));
|
||||
}
|
||||
Poll::Pending => {
|
||||
let Some(stall_timeout) = *this.stall_timeout else {
|
||||
return Poll::Pending;
|
||||
};
|
||||
let timer = this.stall_timer.get_or_insert_with(|| Box::pin(time::sleep(stall_timeout)));
|
||||
if timer.as_mut().poll(cx).is_ready() {
|
||||
record_internode_operation_duration_once(
|
||||
*this.track_internode_metrics,
|
||||
*this.internode_operation,
|
||||
*this.request_started,
|
||||
this.duration_recorded,
|
||||
);
|
||||
record_internode_stall_timeout(*this.track_internode_metrics, *this.internode_operation);
|
||||
record_internode_error(*this.track_internode_metrics, *this.internode_operation);
|
||||
return Poll::Ready(Err(Error::new(
|
||||
io::ErrorKind::TimedOut,
|
||||
"HttpReader stall timeout: no data received before deadline",
|
||||
)));
|
||||
}
|
||||
return Poll::Pending;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl EtagResolvable for HttpReader {
|
||||
fn is_etag_reader(&self) -> bool {
|
||||
false
|
||||
@@ -2012,6 +2188,144 @@ mod tests {
|
||||
handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn http_chunk_reader_handoff_preserves_boundaries_and_eof() {
|
||||
let state = TestState::default();
|
||||
let Some((url, handle)) = start_test_server(state.clone()).await else {
|
||||
return;
|
||||
};
|
||||
let mut reader = HttpChunkReader::new_with_stall_timeout(url, Method::GET, HeaderMap::new(), None, None)
|
||||
.await
|
||||
.expect("reader should open");
|
||||
assert_eq!(reader.consecutive_empty_chunks, 0);
|
||||
let zero = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 0))
|
||||
.await
|
||||
.expect_err("zero chunk bound is invalid");
|
||||
assert_eq!(zero.kind(), io::ErrorKind::InvalidInput);
|
||||
|
||||
let first = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 2))
|
||||
.await
|
||||
.expect("first chunk read should succeed")
|
||||
.expect("first chunk should not be EOF");
|
||||
assert_eq!(first, b"he"[..]);
|
||||
|
||||
let second = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 8))
|
||||
.await
|
||||
.expect("second chunk read should succeed")
|
||||
.expect("second chunk should not be EOF");
|
||||
assert_eq!(second, b"llo"[..]);
|
||||
|
||||
let eof = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 8))
|
||||
.await
|
||||
.expect("EOF should not be an error");
|
||||
assert!(eof.is_none());
|
||||
assert_eq!(state.get_count.load(Ordering::SeqCst), 1);
|
||||
handle.abort();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn http_chunk_reader_rejects_excessive_empty_chunks_on_both_interfaces() {
|
||||
let make_reader = |inner: HttpByteStream| HttpChunkReader {
|
||||
track_internode_metrics: false,
|
||||
internode_operation: None,
|
||||
stall_timeout: None,
|
||||
stall_timer: None,
|
||||
request_started: Instant::now(),
|
||||
duration_recorded: false,
|
||||
consecutive_empty_chunks: 0,
|
||||
inner,
|
||||
current: None,
|
||||
};
|
||||
let empty_chunks_then_data = |empty_chunks| {
|
||||
let items = (0..empty_chunks)
|
||||
.map(|_| Ok(Bytes::new()))
|
||||
.chain(std::iter::once(Ok(Bytes::from_static(b"data"))));
|
||||
make_reader(Box::pin(stream::iter(items)))
|
||||
};
|
||||
|
||||
let mut cx = Context::from_waker(std::task::Waker::noop());
|
||||
let mut chunk_reader = empty_chunks_then_data(MAX_CONSECUTIVE_EMPTY_CHUNKS - 1);
|
||||
let Poll::Ready(Ok(Some(chunk))) = Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 4) else {
|
||||
panic!("data after fewer than the maximum empty chunks should be returned");
|
||||
};
|
||||
assert_eq!(chunk, b"data"[..]);
|
||||
|
||||
let mut async_reader = empty_chunks_then_data(MAX_CONSECUTIVE_EMPTY_CHUNKS - 1);
|
||||
let mut storage = [0; 4];
|
||||
let mut read_buf = ReadBuf::new(&mut storage);
|
||||
let Poll::Ready(Ok(())) = Pin::new(&mut async_reader).poll_read(&mut cx, &mut read_buf) else {
|
||||
panic!("AsyncRead should return data after fewer than the maximum empty chunks");
|
||||
};
|
||||
assert_eq!(read_buf.filled(), b"data");
|
||||
|
||||
let mut chunk_reader = empty_chunks_then_data(MAX_CONSECUTIVE_EMPTY_CHUNKS);
|
||||
let Poll::Ready(Err(err)) = Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 4) else {
|
||||
panic!("excessive empty chunks should fail closed");
|
||||
};
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
let Poll::Ready(Err(err)) = Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 4) else {
|
||||
panic!("empty chunk limit failure should remain sticky");
|
||||
};
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
|
||||
let mut async_reader = empty_chunks_then_data(MAX_CONSECUTIVE_EMPTY_CHUNKS);
|
||||
let mut storage = [0; 4];
|
||||
let mut read_buf = ReadBuf::new(&mut storage);
|
||||
let Poll::Ready(Err(err)) = Pin::new(&mut async_reader).poll_read(&mut cx, &mut read_buf) else {
|
||||
panic!("excessive empty chunks should fail closed through AsyncRead");
|
||||
};
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
let mut read_buf = ReadBuf::new(&mut storage);
|
||||
let Poll::Ready(Err(err)) = Pin::new(&mut async_reader).poll_read(&mut cx, &mut read_buf) else {
|
||||
panic!("AsyncRead empty chunk limit failure should remain sticky");
|
||||
};
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
|
||||
let make_pending_reader = || {
|
||||
let mut empty_chunks = 0;
|
||||
let mut returned_pending = false;
|
||||
let items = stream::poll_fn(move |cx| {
|
||||
if empty_chunks < MAX_CONSECUTIVE_EMPTY_CHUNKS - 1 {
|
||||
empty_chunks += 1;
|
||||
return Poll::Ready(Some(Ok(Bytes::new())));
|
||||
}
|
||||
if !returned_pending {
|
||||
returned_pending = true;
|
||||
cx.waker().wake_by_ref();
|
||||
return Poll::Pending;
|
||||
}
|
||||
if empty_chunks < MAX_CONSECUTIVE_EMPTY_CHUNKS {
|
||||
empty_chunks += 1;
|
||||
return Poll::Ready(Some(Ok(Bytes::new())));
|
||||
}
|
||||
Poll::Ready(None)
|
||||
});
|
||||
make_reader(Box::pin(items))
|
||||
};
|
||||
|
||||
let mut chunk_reader = make_pending_reader();
|
||||
assert!(Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 4).is_pending());
|
||||
let Poll::Ready(Err(err)) = Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 4) else {
|
||||
panic!("the consecutive empty chunk limit must survive Pending");
|
||||
};
|
||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||
|
||||
let items = (0..MAX_CONSECUTIVE_EMPTY_CHUNKS - 1)
|
||||
.map(|_| Ok(Bytes::new()))
|
||||
.chain(std::iter::once(Ok(Bytes::from_static(b"one"))))
|
||||
.chain((0..MAX_CONSECUTIVE_EMPTY_CHUNKS - 1).map(|_| Ok(Bytes::new())))
|
||||
.chain(std::iter::once(Ok(Bytes::from_static(b"two"))));
|
||||
let mut chunk_reader = make_reader(Box::pin(stream::iter(items)));
|
||||
let Poll::Ready(Ok(Some(first))) = Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 3) else {
|
||||
panic!("data should reset the consecutive empty chunk count");
|
||||
};
|
||||
assert_eq!(first, b"one"[..]);
|
||||
let Poll::Ready(Ok(Some(second))) = Pin::new(&mut chunk_reader).poll_read_chunk(&mut cx, 3) else {
|
||||
panic!("empty chunks after data should start a new sequence");
|
||||
};
|
||||
assert_eq!(second, b"two"[..]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn http_reader_records_walk_dir_recv_bytes() {
|
||||
let state = TestState::default();
|
||||
@@ -2347,6 +2661,34 @@ mod tests {
|
||||
handle.abort();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn http_chunk_reader_surfaces_body_error_after_partial_data() {
|
||||
let state = TestState::default();
|
||||
let Some((base_url, handle)) = start_test_server(state).await else {
|
||||
return;
|
||||
};
|
||||
let url = base_url.replace("/stream", "/fail-after-partial");
|
||||
let mut reader = HttpChunkReader::new_with_stall_timeout(url, Method::GET, HeaderMap::new(), None, None)
|
||||
.await
|
||||
.expect("chunk reader should accept the successful response headers");
|
||||
let chunk = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 64))
|
||||
.await
|
||||
.expect("partial response bytes should arrive before the terminal error")
|
||||
.expect("partial response should not be EOF");
|
||||
let err = std::future::poll_fn(|cx| Pin::new(&mut reader).poll_read_chunk(cx, 64))
|
||||
.await
|
||||
.expect_err("terminal body errors must not become clean EOF");
|
||||
|
||||
assert_eq!(chunk, b"partial"[..]);
|
||||
let source = err
|
||||
.get_ref()
|
||||
.and_then(|source| source.downcast_ref::<InternodeHttpError>())
|
||||
.expect("body error should retain internode classification");
|
||||
assert_eq!(source.kind(), InternodeHttpErrorKind::BodyStreamAborted);
|
||||
|
||||
handle.abort();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classify_http_status_marks_retryable_gateway_errors() {
|
||||
let unavailable = classify_http_status(reqwest::StatusCode::SERVICE_UNAVAILABLE);
|
||||
|
||||
Reference in New Issue
Block a user