This is an automated email from the ASF dual-hosted git repository. yiguolei pushed a commit to branch branch-4.2 in repository https://gitbox.apache.org/repos/asf/doris.git
commit b8d6c254b6ac61bff6d6b4cb06ee7cf0ba532188 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]
