This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new 54c1d1ae196 branch-4.1: [fix](s3) Keep the response stream usable when
an error body overflows the read buffer (#68544)
54c1d1ae196 is described below
commit 54c1d1ae1966ac8fca47a455558ec9970b9b75bd
Author: Xin Liao <[email protected]>
AuthorDate: Mon Sep 28 18:54:09 2026 +0800
branch-4.1: [fix](s3) Keep the response stream usable when an error body
overflows the read buffer (#68544)
Pick apache/doris#66557
Adapt the object client changes and test include to branch-4.1's
`be/src/io/fs` layout, preserving its interfaces and retry loop.
Includes all nine response stream unit tests.
Validation: clang-format 16, `build-support/check-format.sh`, and `git
diff --check` passed. Local ASAN BE build and filtered unit test run
were attempted; both were blocked before compilation by the missing
third-party library `libsimdutf.a`. Unit tests were not executed.
---
be/src/io/fs/err_utils.cpp | 9 +-
be/src/io/fs/s3_common.h | 145 ++++++++++++++++++++++-
be/src/io/fs/s3_file_reader.cpp | 10 +-
be/src/io/fs/s3_obj_storage_client.cpp | 14 ++-
be/test/io/fs/s3_response_stream_test.cpp | 189 ++++++++++++++++++++++++++++++
5 files changed, 353 insertions(+), 14 deletions(-)
diff --git a/be/src/io/fs/err_utils.cpp b/be/src/io/fs/err_utils.cpp
index 96ac7b817e9..0a9514605d4 100644
--- a/be/src/io/fs/err_utils.cpp
+++ b/be/src/io/fs/err_utils.cpp
@@ -124,19 +124,22 @@ Status localfs_error(int posix_errno, std::string_view
msg) {
Status s3fs_error(const Aws::S3::S3Error& err, std::string_view msg) {
using namespace Aws::Http;
+ // A failure raised by the client itself carries no request id, and a
dangling
+ // `request_id=` has been read as a request id of the object storage.
+ std::string request_id = err.GetRequestId().empty() ? "<empty>" :
err.GetRequestId().c_str();
switch (err.GetResponseCode()) {
case HttpResponseCode::NOT_FOUND:
return Status::Error<NOT_FOUND, false>("{}: {} {} code=NOT_FOUND,
type={}, request_id={}",
msg, err.GetExceptionName(),
err.GetMessage(),
- err.GetErrorType(),
err.GetRequestId());
+ err.GetErrorType(), request_id);
case HttpResponseCode::FORBIDDEN:
return Status::Error<PERMISSION_DENIED, false>(
"{}: {} {} code=FORBIDDEN, type={}, request_id={}", msg,
err.GetExceptionName(),
- err.GetMessage(), err.GetErrorType(), err.GetRequestId());
+ err.GetMessage(), err.GetErrorType(), request_id);
default:
return Status::Error<ErrorCode::INTERNAL_ERROR, false>(
"{}: {} {} code={} type={}, request_id={}", msg,
err.GetExceptionName(),
- err.GetMessage(), err.GetResponseCode(), err.GetErrorType(),
err.GetRequestId());
+ err.GetMessage(), err.GetResponseCode(), err.GetErrorType(),
request_id);
}
}
diff --git a/be/src/io/fs/s3_common.h b/be/src/io/fs/s3_common.h
index 2f420455227..c5eaae2162a 100644
--- a/be/src/io/fs/s3_common.h
+++ b/be/src/io/fs/s3_common.h
@@ -20,6 +20,11 @@
#include <aws/core/utils/memory/stl/AWSStreamFwd.h>
#include <aws/core/utils/stream/PreallocatedStreamBuf.h>
+#include <algorithm>
+#include <cstring>
+#include <streambuf>
+#include <vector>
+
namespace doris {
// A non-copying iostream.
@@ -34,12 +39,148 @@ public:
std::iostream(this) {}
};
+// The AWS SDK writes the body of every response into the stream built by the
response
+// stream factory of the request, whatever the status of that response is.
Reading an
+// object range straight into the buffer of the caller therefore breaks as
soon as the
+// server answers with an error: the XML body of a `429 SlowDown` is a few
hundred bytes
+// and does not fit into the buffer of a small range read.
`PreallocatedStreamBuf` does not
+// implement `overflow()`, so the stream turns bad, curl aborts the transfer
with
+// `CURLE_WRITE_ERROR`, and the SDK reports an `INTERNAL_FAILURE` named
"Failed to flush
+// response stream" while never recording the status code of the response.
Both the retry
+// strategy of the SDK and the retry of `S3FileReader` key on that status
code, so an error
+// the server asked us to retry ends up cancelling the query instead.
+//
+// This stream buffer writes into the buffer of the caller as long as the body
fits, which
+// is the case for every successful ranged read, and spills the rest into a
buffer of its
+// own. The stream never turns bad, so the SDK reports the real status code
and can parse
+// the error out of the body.
+class S3ResponseStreamBuf final : public std::streambuf {
+public:
+ // Bodies beyond this size are truncated. Only error documents are
expected to overflow
+ // and one is a few hundred bytes, so this leaves them two orders of
magnitude of room
+ // while bounding what a single failing request can hold. Kept small on
purpose: this
+ // buffer is allocated on the transport thread of the SDK, out of the
reach of the memory
+ // tracker of the query, and every concurrent read that fails holds one of
its own.
+ static constexpr size_t MAX_SPILL_SIZE = 64 * 1024;
+
+ S3ResponseStreamBuf(void* buf, size_t nbytes) :
_buf(static_cast<char*>(buf)) {
+ setp(_buf, _buf + nbytes);
+ setg(_buf, _buf, _buf);
+ }
+
+protected:
+ std::streamsize xsputn(const char* s, std::streamsize n) override {
+ if (!_spilled) {
+ if (n <= epptr() - pptr()) {
+ std::memcpy(pptr(), s, n);
+ pbump(static_cast<int>(n));
+ return n;
+ }
+ _spill_over();
+ }
+ // Saturating on its own: the spill is clamped when it is filled from
the buffer of
+ // the caller, and this must not underflow into an unbounded write if
it ever is not.
+ auto room = _spill.size() < MAX_SPILL_SIZE ? MAX_SPILL_SIZE -
_spill.size() : 0;
+ auto writable = std::min(static_cast<size_t>(n), room);
+ _spill.insert(_spill.end(), s, s + writable);
+ // Always report the whole write as consumed. A short write is what
makes curl
+ // abort the transfer and lose the status code of the response.
+ return n;
+ }
+
+ int_type overflow(int_type ch) override {
+ if (traits_type::eq_int_type(ch, traits_type::eof())) {
+ return traits_type::not_eof(ch);
+ }
+ auto c = traits_type::to_char_type(ch);
+ xsputn(&c, 1);
+ return ch;
+ }
+
+ int_type underflow() override {
+ _reset_get_area(_read_pos());
+ if (gptr() == egptr()) {
+ return traits_type::eof();
+ }
+ return traits_type::to_int_type(*gptr());
+ }
+
+ pos_type seekoff(off_type off, std::ios_base::seekdir dir,
+ std::ios_base::openmode which) override {
+ auto size = static_cast<off_type>(_written());
+ if ((which & std::ios_base::out) && !(which & std::ios_base::in)) {
+ // The SDK only asks for the write position, to tell an empty body
apart from a
+ // body it has to parse. Moving the write pointer is not supported.
+ return dir == std::ios_base::cur && off == 0 ? pos_type(size) :
pos_type(off_type(-1));
+ }
+ // A seek asking for both areas at once, which is what the default
argument of
+ // `pubseekoff()` and `pubseekpos()` does, is served as a seek of the
read area. The
+ // write area is append only, so there is nothing to move there.
+ off_type pos = off;
+ if (dir == std::ios_base::cur) {
+ pos += static_cast<off_type>(_read_pos());
+ } else if (dir == std::ios_base::end) {
+ pos += size;
+ }
+ if (pos < 0 || pos > size) {
+ return pos_type(off_type(-1));
+ }
+ _reset_get_area(static_cast<size_t>(pos));
+ return pos_type(pos);
+ }
+
+ pos_type seekpos(pos_type pos, std::ios_base::openmode which) override {
+ return seekoff(pos, std::ios_base::beg, which);
+ }
+
+private:
+ // Moves what has been written so far into the spill buffer, so that the
body stays
+ // contiguous and the SDK can parse the error out of it. Truncated right
here: the buffer
+ // of the caller is the size of the range that was asked for,
`remote_storage_read_buffer_mb`
+ // of it for a prefetched read and the whole file for a download, so it
can be far larger
+ // than the bound of the spill. Starting the spill beyond its own bound
would leave no room
+ // for the truncation to ever apply and let a server answering a ranged
read with the whole
+ // object be buffered in full.
+ void _spill_over() {
+ auto kept = std::min(static_cast<size_t>(pptr() - _buf),
MAX_SPILL_SIZE);
+ _spill.assign(_buf, _buf + kept);
+ setp(nullptr, nullptr);
+ _spilled = true;
+ }
+
+ // Bytes of the body held by this buffer, truncation excluded.
+ size_t _written() const { return _spilled ? _spill.size() : pptr() - _buf;
}
+
+ // Both areas start at the same logical offset, so the read position
survives a spill.
+ size_t _read_pos() const { return gptr() - eback(); }
+
+ void _reset_get_area(size_t pos) {
+ char* begin = _spilled ? _spill.data() : _buf;
+ auto size = _written();
+ pos = std::min(pos, size);
+ setg(begin, begin + pos, begin + size);
+ }
+
+ char* _buf;
+ std::vector<char> _spill;
+ bool _spilled = false;
+};
+
+class S3ResponseStream final : public std::iostream {
+public:
+ S3ResponseStream(void* buf, size_t nbytes) : std::iostream(&_buf),
_buf(buf, nbytes) {}
+
+private:
+ S3ResponseStreamBuf _buf;
+};
+
// By default, the AWS SDK reads object data into an auto-growing StringStream.
-// To avoid copies, read directly into our preallocated buffer instead.
+// To avoid copies, read the body directly into our preallocated buffer
instead, and keep
+// only what does not fit, which is an error document, in a buffer of the
stream itself.
// See https://github.com/aws/aws-sdk-cpp/issues/64 for an alternative but
// functionally similar recipe.
inline Aws::IOStreamFactory AwsWriteableStreamFactory(void* buf, int64_t
nbytes) {
- return [=]() { return Aws::New<StringViewStream>("", buf, nbytes); };
+ return [=]() { return Aws::New<S3ResponseStream>("", buf,
static_cast<size_t>(nbytes)); };
}
} // namespace doris
diff --git a/be/src/io/fs/s3_file_reader.cpp b/be/src/io/fs/s3_file_reader.cpp
index af8dde36d2d..dbbfd8240e6 100644
--- a/be/src/io/fs/s3_file_reader.cpp
+++ b/be/src/io/fs/s3_file_reader.cpp
@@ -160,6 +160,7 @@ Status S3FileReader::read_at_impl(size_t offset, Slice
result, size_t* bytes_rea
SCOPED_RAW_TIMER(&_s3_stats.total_get_request_time_ns);
int total_sleep_time = 0;
+ Status last_error;
while (retry_count <= max_retries) {
*bytes_read = 0;
s3_file_reader_read_counter << 1;
@@ -172,6 +173,7 @@ Status S3FileReader::read_at_impl(size_t offset, Slice
result, size_t* bytes_rea
if (resp.http_code ==
static_cast<int>(Aws::Http::HttpResponseCode::TOO_MANY_REQUESTS)) {
s3_file_reader_too_many_request_counter << 1;
+ last_error = Status(resp.status.code,
std::move(resp.status.msg));
retry_count++;
int wait_time = std::min(base_wait_time * (1 << retry_count),
max_wait_time); // Exponential backoff
@@ -182,8 +184,7 @@ Status S3FileReader::read_at_impl(size_t offset, Slice
result, size_t* bytes_rea
continue;
} else {
// Handle other errors
- return std::move(Status(resp.status.code,
std::move(resp.status.msg))
- .append("failed to read"));
+ return {resp.status.code, std::move(resp.status.msg)};
}
}
if (*bytes_read != bytes_req) {
@@ -206,8 +207,9 @@ Status S3FileReader::read_at_impl(size_t offset, Slice
result, size_t* bytes_rea
}
std::string msg = fmt::format(
"failed to get object, path={} offset={} bytes_req={}
bytes_read={} file_size={} "
- "tries={}",
- _path.native(), offset, bytes_req, *bytes_read, _file_size,
(max_retries + 1));
+ "tries={}, last error: [{}]",
+ _path.native(), offset, bytes_req, *bytes_read, _file_size,
(max_retries + 1),
+ last_error.msg());
LOG(WARNING) << msg;
return Status::InternalError(msg);
}
diff --git a/be/src/io/fs/s3_obj_storage_client.cpp
b/be/src/io/fs/s3_obj_storage_client.cpp
index 0c0b0370f80..f553de8ca48 100644
--- a/be/src/io/fs/s3_obj_storage_client.cpp
+++ b/be/src/io/fs/s3_obj_storage_client.cpp
@@ -310,17 +310,21 @@ ObjectStorageResponse
S3ObjStorageClient::get_object(const ObjectStoragePathOpti
if (!outcome.IsSuccess()) {
record_s3_request_failed(outcome.GetError());
return {convert_to_obj_response(s3fs_error(
- outcome.GetError(), fmt::format("failed to read from
{}", opts.key))),
+ outcome.GetError(),
+ fmt::format("failed to get object: bucket={} object={}
offset={} size={}",
+ opts.bucket, opts.key, offset,
bytes_read))),
static_cast<int>(outcome.GetError().GetResponseCode()),
outcome.GetError().GetRequestId()};
}
*size_return = outcome.GetResult().GetContentLength();
- // case for incomplete read
+ // Short read, or a server or a proxy answering a ranged read with the
whole object.
SYNC_POINT_CALLBACK("s3_obj_storage_client::get_object", size_return);
if (*size_return != bytes_read) {
- return {convert_to_obj_response(Status::InternalError(
- "failed to read from {}(bytes read: {}, bytes req: {}),
request_id: {}", opts.key,
- *size_return, bytes_read,
outcome.GetResult().GetRequestId()))};
+ return {convert_to_obj_response(
+ Status::InternalError("incomplete read from bucket={}
object={} offset={}, expect "
+ "{}, got {}, request_id={}",
+ opts.bucket, opts.key, offset,
bytes_read, *size_return,
+ outcome.GetResult().GetRequestId()))};
}
return ObjectStorageResponse::OK();
}
diff --git a/be/test/io/fs/s3_response_stream_test.cpp
b/be/test/io/fs/s3_response_stream_test.cpp
new file mode 100644
index 00000000000..65b246f7c1e
--- /dev/null
+++ b/be/test/io/fs/s3_response_stream_test.cpp
@@ -0,0 +1,189 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include <gtest/gtest.h>
+
+#include <sstream>
+#include <string>
+#include <vector>
+
+#include "io/fs/s3_common.h"
+
+namespace doris {
+
+namespace {
+
+// What the SDK does with the body of a response it has to build an error from.
+std::string drain(std::iostream& stream) {
+ std::stringstream out;
+ out << stream.rdbuf();
+ return out.str();
+}
+
+// The XML body a MinIO answers a throttled ranged read with, shortened.
+constexpr char SLOW_DOWN_BODY[] =
+ R"(<?xml version="1.0"
encoding="UTF-8"?><Error><Code>SlowDown</Code><Message>Please )"
+ R"(reduce your request
rate.</Message><Key>data/packed_file/2666/x.bin</Key></Error>)";
+
+} // namespace
+
+// A body of the requested size lands in the buffer of the caller, without a
copy.
+TEST(S3ResponseStreamTest, BodyFits) {
+ std::string body(64, 'a');
+ std::vector<char> buffer(body.size());
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ stream.write(body.data(), body.size());
+ stream.flush();
+
+ EXPECT_FALSE(stream.fail());
+ EXPECT_EQ(body, std::string(buffer.data(), buffer.size()));
+ EXPECT_EQ(static_cast<std::streampos>(body.size()), stream.tellp());
+ EXPECT_EQ(body, drain(stream));
+}
+
+// An error body larger than the range of the read leaves the stream usable,
which is what
+// keeps curl from aborting the transfer and the SDK from losing the status
code.
+TEST(S3ResponseStreamTest, ErrorBodyOverflowsInOneWrite) {
+ std::string body(SLOW_DOWN_BODY);
+ // A read of the footer of a packed file is far smaller than the error
document.
+ std::vector<char> buffer(12);
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ stream.write(body.data(), body.size());
+ stream.flush();
+
+ EXPECT_FALSE(stream.fail());
+ EXPECT_EQ(static_cast<std::streampos>(body.size()), stream.tellp());
+ EXPECT_EQ(body, drain(stream));
+}
+
+// curl hands the body over in chunks, so the overflow can happen in the
middle of one.
+TEST(S3ResponseStreamTest, ErrorBodyOverflowsAcrossWrites) {
+ std::string body(SLOW_DOWN_BODY);
+ std::vector<char> buffer(16);
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ size_t chunk = 7;
+ for (size_t pos = 0; pos < body.size(); pos += chunk) {
+ stream.write(body.data() + pos, std::min(chunk, body.size() - pos));
+ }
+ stream.flush();
+
+ EXPECT_FALSE(stream.fail());
+ EXPECT_EQ(static_cast<std::streampos>(body.size()), stream.tellp());
+ // The bytes written before the overflow are kept, so the body stays
contiguous.
+ EXPECT_EQ(body, drain(stream));
+}
+
+// A body written one character at a time goes through overflow() instead of
xsputn().
+TEST(S3ResponseStreamTest, ErrorBodyOverflowsCharByChar) {
+ std::string body(SLOW_DOWN_BODY);
+ std::vector<char> buffer(4);
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ for (char c : body) {
+ stream.put(c);
+ }
+ stream.flush();
+
+ EXPECT_FALSE(stream.fail());
+ EXPECT_EQ(body, drain(stream));
+}
+
+// A server answering a ranged read with the whole object must not blow up the
memory of the
+// backend. The body is truncated, the stream stays good and the read is
rejected later on by
+// the length check of the caller.
+TEST(S3ResponseStreamTest, OversizedBodyIsTruncated) {
+ std::string body(S3ResponseStreamBuf::MAX_SPILL_SIZE + 4096, 'x');
+ std::vector<char> buffer(8);
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ stream.write(body.data(), body.size());
+ stream.flush();
+
+ EXPECT_FALSE(stream.fail());
+
EXPECT_EQ(static_cast<std::streampos>(S3ResponseStreamBuf::MAX_SPILL_SIZE),
stream.tellp());
+ EXPECT_EQ(S3ResponseStreamBuf::MAX_SPILL_SIZE, drain(stream).size());
+}
+
+// The same, with a buffer of the caller larger than the bound of the spill: a
prefetched
+// read asks for `remote_storage_read_buffer_mb` at a time and a download for
the whole file,
+// so the bytes moved out of that buffer on the overflow have to be truncated
as well.
+TEST(S3ResponseStreamTest, OversizedBodyIsTruncatedWithLargeBuffer) {
+ std::vector<char> buffer(2 * S3ResponseStreamBuf::MAX_SPILL_SIZE);
+ std::string body(4 * S3ResponseStreamBuf::MAX_SPILL_SIZE, 'x');
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ // curl hands the body over in chunks of `CURL_MAX_WRITE_SIZE`, so the
buffer of the caller
+ // is filled before a write overflows it.
+ size_t chunk = 16384;
+ for (size_t pos = 0; pos < body.size(); pos += chunk) {
+ stream.write(body.data() + pos, std::min(chunk, body.size() - pos));
+ }
+ stream.flush();
+
+ EXPECT_FALSE(stream.fail());
+
EXPECT_EQ(static_cast<std::streampos>(S3ResponseStreamBuf::MAX_SPILL_SIZE),
stream.tellp());
+ EXPECT_EQ(S3ResponseStreamBuf::MAX_SPILL_SIZE, drain(stream).size());
+}
+
+// A truncated body is still a body the SDK rewinds and reads to its end.
+TEST(S3ResponseStreamTest, SeekTruncatedBody) {
+ std::vector<char> buffer(2 * S3ResponseStreamBuf::MAX_SPILL_SIZE);
+ std::string body(4 * S3ResponseStreamBuf::MAX_SPILL_SIZE, 'x');
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ stream.write(body.data(), body.size());
+
+ EXPECT_EQ(S3ResponseStreamBuf::MAX_SPILL_SIZE, drain(stream).size());
+ stream.clear();
+ EXPECT_EQ(std::streampos(0), stream.seekg(0).tellg());
+ EXPECT_EQ(S3ResponseStreamBuf::MAX_SPILL_SIZE, drain(stream).size());
+ // Past the end of what has been kept.
+ stream.clear();
+ EXPECT_TRUE(stream.seekg(S3ResponseStreamBuf::MAX_SPILL_SIZE + 1).fail());
+}
+
+// The SDK rewinds the body before parsing an error out of it.
+TEST(S3ResponseStreamTest, SeekBackAndForth) {
+ std::string body(SLOW_DOWN_BODY);
+ std::vector<char> buffer(12);
+
+ S3ResponseStream stream(buffer.data(), buffer.size());
+ stream.write(body.data(), body.size());
+
+ EXPECT_EQ(body, drain(stream));
+ stream.clear();
+ stream.seekg(0);
+ EXPECT_EQ(body, drain(stream));
+
+ stream.clear();
+ stream.seekg(2);
+ EXPECT_EQ(body.substr(2), drain(stream));
+}
+
+// An empty body is what tells the SDK to build the error out of the status
code alone.
+TEST(S3ResponseStreamTest, EmptyBody) {
+ std::vector<char> buffer(16);
+ S3ResponseStream stream(buffer.data(), buffer.size());
+
+ EXPECT_EQ(std::streampos(0), stream.tellp());
+ EXPECT_TRUE(drain(stream).empty());
+}
+
+} // namespace doris
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]