Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,10 @@ set(util_source_files
src/util/byte_span_inserters.hpp
src/util/byte_span_packed_int_constants.hpp

src/util/byte_range_fwd.hpp
src/util/byte_range.hpp
src/util/byte_range.cpp

src/util/command_line_helpers_fwd.hpp
src/util/command_line_helpers.hpp
src/util/command_line_helpers.cpp
Expand Down
2 changes: 1 addition & 1 deletion extra/mysql_protocol/mysql/harness/stdx/ranges.h
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@

namespace stdx::ranges {

// TODO: change the content of this this file to "using
// TODO: change the content of this file to "using
// std::ranges::views::enumerate" when switching to clang-23

/**
Expand Down
2 changes: 1 addition & 1 deletion mtr/binlog_streaming/t/checkpointing.test
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ DROP TABLE t1;

# creating data directory, configuration file, etc.

# reducing log level as this this generates a big number of binlog events
# reducing log level as this generates a big number of binlog events
# that do not need to be logged in details
--let $binsrv_log_level = info
--let $binsrv_connect_timeout = 20
Expand Down
6 changes: 4 additions & 2 deletions src/binsrv/basic_storage_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
#include <string>
#include <string_view>

#include "util/byte_range_fwd.hpp"
#include "util/byte_span_fwd.hpp"
#include "util/exception_location_helpers.hpp"

Expand All @@ -33,8 +34,9 @@ basic_storage_backend::list_objects() {
}

[[nodiscard]] std::string
basic_storage_backend::get_object(std::string_view name) {
return do_get_object(name);
basic_storage_backend::get_object(std::string_view name,
const util::byte_range &range) {
return do_get_object(name, range);
}

void basic_storage_backend::put_object(std::string_view name,
Expand Down
12 changes: 10 additions & 2 deletions src/binsrv/basic_storage_backend.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,12 +22,17 @@
#include <string>
#include <string_view>

#include "util/byte_range.hpp"
#include "util/byte_span_fwd.hpp"
#include "util/common_optional_types.hpp"

namespace binsrv {

class basic_storage_backend {
public:
// 256 MB
static constexpr std::size_t max_memory_object_size{256UZ << 20UZ};

basic_storage_backend() = default;
basic_storage_backend(const basic_storage_backend &) = delete;
basic_storage_backend(basic_storage_backend &&) noexcept = delete;
Expand All @@ -37,7 +42,9 @@ class basic_storage_backend {
virtual ~basic_storage_backend() = default;

[[nodiscard]] storage_object_name_container list_objects();
[[nodiscard]] std::string get_object(std::string_view name);
[[nodiscard]] std::string
get_object(std::string_view name,
const util::byte_range &range = util::byte_range{});
// 'put_object' is an atomic overwrite: a concurrent / post-crash
// reader either sees the previous bytes in full or the new bytes in
// full, never a partial mix.
Expand Down Expand Up @@ -71,7 +78,8 @@ class basic_storage_backend {
bool stream_open_{false};

[[nodiscard]] virtual storage_object_name_container do_list_objects() = 0;
[[nodiscard]] virtual std::string do_get_object(std::string_view name) = 0;
[[nodiscard]] virtual std::string
do_get_object(std::string_view name, const util::byte_range &range) = 0;
virtual void do_put_object(std::string_view name,
util::const_byte_span content) = 0;
virtual void do_resize_object(std::string_view name,
Expand Down
12 changes: 7 additions & 5 deletions src/binsrv/filesystem_storage_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@

#include "binsrv/storage_config.hpp"

#include "util/byte_range_fwd.hpp"
#include "util/byte_span.hpp"
#include "util/exception_location_helpers.hpp"
#include "util/file_operations_helpers.hpp"
Expand Down Expand Up @@ -115,10 +116,11 @@ filesystem_storage_backend::do_list_objects() {
}

[[nodiscard]] std::string
filesystem_storage_backend::do_get_object(std::string_view name) {
filesystem_storage_backend::do_get_object(std::string_view name,
const util::byte_range &range) {
const auto object_path{get_object_path(name)};
return util::read_file_content(object_path, max_memory_object_size,
"underlying object file");
return util::read_file_content("underlying object file", object_path,
max_memory_object_size, range);
}

void filesystem_storage_backend::do_put_object(std::string_view name,
Expand All @@ -138,8 +140,8 @@ void filesystem_storage_backend::do_put_object(std::string_view name,
auto tmp_object_path = object_path;
tmp_object_path += tmp_storage_object_suffix;

util::write_file_content(tmp_object_path, util::as_string_view(content),
"underlying tmp object file");
util::write_file_content("underlying tmp object file", tmp_object_path,
util::as_string_view(content));
// make the tmp file's content durable before the rename swaps it
util::fsync(tmp_object_path);

Expand Down
7 changes: 4 additions & 3 deletions src/binsrv/filesystem_storage_backend.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,13 @@
#include "binsrv/basic_storage_backend.hpp" // IWYU pragma: export
#include "binsrv/storage_config_fwd.hpp"

#include "util/byte_range_fwd.hpp"

namespace binsrv {

class [[nodiscard]] filesystem_storage_backend final
: public basic_storage_backend {
public:
static constexpr std::size_t max_memory_object_size{1048576U};

static constexpr std::string_view uri_schema{"file"};

explicit filesystem_storage_backend(const storage_config &config);
Expand All @@ -44,7 +44,8 @@ class [[nodiscard]] filesystem_storage_backend final

[[nodiscard]] storage_object_name_container do_list_objects() override;

[[nodiscard]] std::string do_get_object(std::string_view name) override;
[[nodiscard]] std::string
do_get_object(std::string_view name, const util::byte_range &range) override;
void do_put_object(std::string_view name,
util::const_byte_span content) override;
void do_resize_object(std::string_view name, std::uint64_t new_size) override;
Expand Down
4 changes: 2 additions & 2 deletions src/binsrv/keyring_record_collection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -35,8 +35,8 @@ namespace binsrv {
keyring_record_collection::keyring_record_collection(std::string_view file_name)
: impl_{} {
static constexpr std::size_t max_file_size{1048576U};
const auto data = util::read_file_content(file_name, max_file_size,
"keyring record collection file");
const auto data = util::read_file_content("keyring record collection file",
file_name, max_file_size);
auto json_value = boost::json::parse(data);
util::nv_tuple_from_json(json_value, impl_);

Expand Down
2 changes: 1 addition & 1 deletion src/binsrv/main_config.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ main_config::main_config(std::string_view file_name) {
static constexpr std::size_t max_file_size{1048576U};

const auto file_content =
util::read_file_content(file_name, max_file_size, "configuration file");
util::read_file_content("configuration file", file_name, max_file_size);
if (file_content.empty()) {
util::exception_location().raise<std::out_of_range>(
"configuration file is empty");
Expand Down
89 changes: 69 additions & 20 deletions src/binsrv/s3_storage_backend.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@
#include "binsrv/s3_error_helpers_private.hpp"
#include "binsrv/storage_config.hpp"

#include "util/byte_range.hpp"
#include "util/byte_span.hpp"
#include "util/exception_location_helpers.hpp"

Expand Down Expand Up @@ -160,12 +161,14 @@ class s3_storage_backend::aws_context : private aws_context_base {

[[nodiscard]] std::string get_bucket_region(const std::string &bucket) const;

[[nodiscard]] std::string
get_object_into_string(const qualified_object_path &source) const;
[[nodiscard]] std::string get_object_into_string(
const qualified_object_path &source,
const util::byte_range &range = util::byte_range{}) const;

void
get_object_into_file(const qualified_object_path &source,
const std::filesystem::path &content_file_path) const;
void get_object_into_file(
const qualified_object_path &source,
const std::filesystem::path &content_file_path,
const util::byte_range &range = util::byte_range{}) const;

void put_object_from_stream(const qualified_object_path &dest,
std::iostream &content_stream) const;
Expand Down Expand Up @@ -193,7 +196,8 @@ class s3_storage_backend::aws_context : private aws_context_base {

void get_object_internal(const qualified_object_path &source,
const stream_factory_type &stream_factory,
const stream_handler_type &stream_handler) const;
const stream_handler_type &stream_handler,
const util::byte_range &range) const;

using list_object_container = Aws::Vector<Aws::S3Crt::Model::Object>;
static void
Expand Down Expand Up @@ -268,15 +272,36 @@ s3_storage_backend::aws_context::aws_context(

[[nodiscard]] std::string
s3_storage_backend::aws_context::get_object_into_string(
const qualified_object_path &source) const {
const qualified_object_path &source, const util::byte_range &range) const {
if (range.is_empty()) {
return {};
}

if (range.has_length()) {
if (range.get_length() > max_memory_object_size) {
util::exception_location().raise<std::out_of_range>(
"The requested S3 object range is too large to be loaded in memory");
}
}
std::string content;
auto stream_handler{[&content](std::size_t content_length,
std::iostream &content_stream) {
auto stream_handler{[&content, &range](std::size_t content_length,
std::iostream &content_stream) {
// TODO: check object length in advance before calling GetObject
// (with HeadObject, for instance)
if (content_length > max_memory_object_size) {
util::exception_location().raise<std::out_of_range>(
"S3 object is too large to be loaded in memory");
// alternatively, set `bytes=0-<max_memory_object_size - 1>` byte
// range in the request and this operation will return up to
// `max_memory_object_size` bytes.
if (range.has_length()) {
if (content_length != range.get_length()) {
util::exception_location().raise<std::out_of_range>(
"The requested S3 object range does not match the length of the "
"received memory content");
}
} else {
if (content_length > max_memory_object_size) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

isn't stream_handler called after receiving the object from s3? So it is already in the memory.

@percona-ysorokin percona-ysorokin Sep 23, 2026 •

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, this is a compromise we had to make here. Unfortunately, there is no way in S3 API to "read no more than N bytes from an S3 object" in a single request. So, we either need to send another request (HeadObject) before this call, or simply rely on the fact that we won't be abusing this function.
See the TODO item above.
I don't like this either to be honest.

util::exception_location().raise<std::out_of_range>(
"S3 object is too large to be loaded in memory");
Comment on lines +294 to +303
}
}

content.resize(content_length);
Expand All @@ -288,14 +313,24 @@ s3_storage_backend::aws_context::get_object_into_string(
assert(content_stream.gcount() ==
static_cast<std::streamsize>(content_length));
}};
get_object_internal(source, {}, stream_handler);
get_object_internal(source, {}, stream_handler, range);

return content;
}

void s3_storage_backend::aws_context::get_object_into_file(
const qualified_object_path &source,
const std::filesystem::path &content_file_path) const {
const std::filesystem::path &content_file_path,
const util::byte_range &range) const {

if (range.is_empty()) {
// we need to create an empty file in this case
const std::ofstream empty_file{content_file_path,
std::ios_base::out | std::ios_base::binary |
std::ios_base::trunc};
return;
}

auto stream_factory{[&content_file_path]() -> std::iostream * {
return Aws::New<std::fstream>(
"GetObjectStreamFactoryAllocationTag", content_file_path,
Expand All @@ -304,8 +339,15 @@ void s3_storage_backend::aws_context::get_object_into_file(
}};
std::size_t response_content_length{};
auto stream_handler{
[&response_content_length](std::size_t content_length,
std::iostream &content_stream) {
[&response_content_length, &range](std::size_t content_length,
std::iostream &content_stream) {
if (range.has_length()) {
if (content_length != range.get_length()) {
util::exception_location().raise<std::out_of_range>(
"The requested S3 object range does not match the length of "
"the received file content");
}
}
content_stream.seekg(0, std::ios_base::end);

const auto end_position{
Expand All @@ -317,7 +359,7 @@ void s3_storage_backend::aws_context::get_object_into_file(
response_content_length = content_length;
}};

get_object_internal(source, stream_factory, stream_handler);
get_object_internal(source, stream_factory, stream_handler, range);
assert(std::filesystem::file_size(content_file_path) ==
response_content_length);
}
Expand Down Expand Up @@ -448,14 +490,20 @@ s3_storage_backend::aws_context::list_objects(
void s3_storage_backend::aws_context::get_object_internal(
const qualified_object_path &source,
const stream_factory_type &stream_factory,
const stream_handler_type &stream_handler) const {
const stream_handler_type &stream_handler,
const util::byte_range &range) const {
assert(!range.is_empty());
Aws::S3Crt::Model::GetObjectRequest get_object_request;
if (stream_factory) {
get_object_request.SetResponseStreamFactory(stream_factory);
}
get_object_request.SetBucket(source.bucket);
get_object_request.SetKey(source.object_path.generic_string());

if (!range.is_full()) {
get_object_request.SetRange("bytes=" + range.to_string());
}

const auto get_object_outcome{client_->GetObject(get_object_request)};

if (!get_object_outcome.IsSuccess()) {
Expand Down Expand Up @@ -700,9 +748,10 @@ s3_storage_backend::do_list_objects() {
}

[[nodiscard]] std::string
s3_storage_backend::do_get_object(std::string_view name) {
s3_storage_backend::do_get_object(std::string_view name,
const util::byte_range &range) {
return impl_->get_object_into_string(
{.bucket = bucket_, .object_path = get_object_path(name)});
{.bucket = bucket_, .object_path = get_object_path(name)}, range);
}

void s3_storage_backend::do_put_object(std::string_view name,
Expand Down
5 changes: 2 additions & 3 deletions src/binsrv/s3_storage_backend.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,6 @@ namespace binsrv {

class [[nodiscard]] s3_storage_backend final : public basic_storage_backend {
public:
static constexpr std::size_t max_memory_object_size{1048576U};

static constexpr std::string_view original_uri_schema{"s3"};

explicit s3_storage_backend(const storage_config &config);
Expand Down Expand Up @@ -70,7 +68,8 @@ class [[nodiscard]] s3_storage_backend final : public basic_storage_backend {

[[nodiscard]] storage_object_name_container do_list_objects() override;

[[nodiscard]] std::string do_get_object(std::string_view name) override;
[[nodiscard]] std::string
do_get_object(std::string_view name, const util::byte_range &range) override;
void do_put_object(std::string_view name,
util::const_byte_span content) override;
void do_resize_object(std::string_view name, std::uint64_t new_size) override;
Expand Down
10 changes: 5 additions & 5 deletions src/binsrv/storage_core.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -235,7 +235,7 @@ void storage_core::set_purged_gtids(const gtids::gtid_set &purged_gtids) {
}

[[nodiscard]] std::string storage_core::get_backend_description() const {
// no mutex protection needed as this this method calls a const
// no mutex protection needed as this method calls a const
// method on an instance of basic_storage_backend that reads only data
// that was set only once during construction
return backend_->get_description();
Expand Down Expand Up @@ -425,29 +425,29 @@ storage_core::purge_binlogs(const events::composite_binlog_name &target) {

[[nodiscard]] std::string storage_core::get_binlog_uri(
const events::composite_binlog_name &binlog_name) const {
// no mutex protection needed as this this method calls a const
// no mutex protection needed as this method calls a const
// method on an instance of basic_storage_backend that reads only data
// that was set only once during construction
return backend_->get_object_uri(binlog_name.str());
}

[[nodiscard]] std::string storage_core::get_keyring_description() const {
// no mutex protection needed as this this method calls a const
// no mutex protection needed as this method calls a const
// method on an immutable keyring instance
return is_keyring_initialized() ? keyring_->get_description()
: "keyring is not initialized";
}

[[nodiscard]] std::string storage_core::get_active_kek_description() const {
// no mutex protection needed as this this method calls a chain of const
// no mutex protection needed as this method calls a chain of const
// methods on an immutable keyring instance
return has_active_kek() ? keyring_->get_key(active_kek_id_).get_description()
: "active KEK is not set";
}

[[nodiscard]] std::string
storage_core::get_encryption_format_description() const {
// no mutex protection needed as this this method reads data
// no mutex protection needed as this method reads data
// set only once during construction
return encryption_format_.has_value()
? std::string{to_string_view(*encryption_format_)}
Expand Down
Loading
Loading