From 6e2620af0dfa3809ab5d423eb659bd42063b4264 Mon Sep 17 00:00:00 2001 From: Yura Sorokin Date: Wed, 23 Sep 2026 11:27:27 +0200 Subject: [PATCH] PBS-6 feature: Rework storage layer abstractions to support simultaneous reads and writes (part 3) https://perconadev.atlassian.net/browse/PBS-6 Implemented the most straightforward locking strategy for 'binsrv::storage_core'. * Added an instance of the 'std::shared_mutex' to the 'binsrv::storage_core' class. * Public const methods are now protected with 'std::shared_lock'. * Public mutable methods are now protected with 'std::unique_lock'. * For public methods that are called from other public methods introduced its 'xxx_unsafe()' version to avoid double locking. * Removed 'noexcept' qualifier from methods that create 'std::unique_lock' / 'std::shared_lock'. * Public const methods that return data that is initialized only once during object construction left without locks. * Some methods that used to return references to internal objects now return by value to achieve thread-safety guarantees. --- src/binsrv/storage.cpp | 16 +++--- src/binsrv/storage.hpp | 15 +++--- src/binsrv/storage_core.cpp | 62 ++++++++++++++++------ src/binsrv/storage_core.hpp | 100 ++++++++++++++++++++++++------------ 4 files changed, 126 insertions(+), 67 deletions(-) diff --git a/src/binsrv/storage.cpp b/src/binsrv/storage.cpp index 634bb3d..518c036 100644 --- a/src/binsrv/storage.cpp +++ b/src/binsrv/storage.cpp @@ -60,8 +60,7 @@ storage::~storage() { } } -[[nodiscard]] const gtids::gtid_set & -storage::get_purged_gtids() const noexcept { +[[nodiscard]] gtids::gtid_set storage::get_purged_gtids() const { return core_->get_purged_gtids(); } @@ -82,14 +81,11 @@ storage::get_replication_mode() const noexcept { return core_->is_in_gtid_replication_mode(); } -[[nodiscard]] const binlog_record_container & -storage::get_binlog_records() const noexcept { +[[nodiscard]] binlog_record_container storage::get_binlog_records() const { return core_->get_binlog_records(); } -[[nodiscard]] bool storage::is_empty() const noexcept { - return core_->is_empty(); -} +[[nodiscard]] bool storage::is_empty() const { return core_->is_empty(); } [[nodiscard]] gtids::gtid_set storage::get_gtids() const { return core_->get_gtids(); @@ -100,7 +96,7 @@ storage::get_current_binlog_name() const { return core_->get_current_binlog_name(); } -[[nodiscard]] bool storage::is_binlog_open() const noexcept { +[[nodiscard]] bool storage::is_binlog_open() const { return core_->is_binlog_open(); } @@ -109,7 +105,7 @@ storage::open_binlog(const events::composite_binlog_name &binlog_name) { auto result{core_->open_binlog(binlog_name)}; if (result != open_binlog_status::created) { - ready_to_flush_last_sequence_number_ = core_->last_sequence_number(); + ready_to_flush_last_sequence_number_ = core_->get_last_sequence_number(); incomplete_transaction_last_sequence_number_ = ready_to_flush_last_sequence_number_; } else { @@ -250,7 +246,7 @@ void storage::update_last_checkpoint_info() { } } -[[nodiscard]] std::uint64_t storage::get_flushed_position() const noexcept { +[[nodiscard]] std::uint64_t storage::get_flushed_position() const { return core_->get_flushed_position(); } diff --git a/src/binsrv/storage.hpp b/src/binsrv/storage.hpp index ede931a..da6329f 100644 --- a/src/binsrv/storage.hpp +++ b/src/binsrv/storage.hpp @@ -63,7 +63,7 @@ class [[nodiscard]] storage { ~storage(); - [[nodiscard]] const gtids::gtid_set &get_purged_gtids() const noexcept; + [[nodiscard]] gtids::gtid_set get_purged_gtids() const; void set_purged_gtids(const gtids::gtid_set &purged_gtids); [[nodiscard]] std::string get_backend_description() const; @@ -71,12 +71,11 @@ class [[nodiscard]] storage { [[nodiscard]] replication_mode_type get_replication_mode() const noexcept; [[nodiscard]] bool is_in_gtid_replication_mode() const noexcept; - [[nodiscard]] const binlog_record_container & - get_binlog_records() const noexcept; - [[nodiscard]] bool is_empty() const noexcept; + [[nodiscard]] binlog_record_container get_binlog_records() const; + [[nodiscard]] bool is_empty() const; [[nodiscard]] events::composite_binlog_name get_current_binlog_name() const; - [[nodiscard]] std::uint64_t get_current_position() const noexcept { + [[nodiscard]] std::uint64_t get_current_position() const { return get_flushed_position() + std::size(event_buffer_); } @@ -87,7 +86,7 @@ class [[nodiscard]] storage { return incomplete_transaction_last_sequence_number_; } - [[nodiscard]] bool is_binlog_open() const noexcept; + [[nodiscard]] bool is_binlog_open() const; [[nodiscard]] open_binlog_status open_binlog(const events::composite_binlog_name &binlog_name); @@ -159,8 +158,8 @@ class [[nodiscard]] storage { [[nodiscard]] bool has_event_data_to_flush() const noexcept { return last_transaction_boundary_position_in_event_buffer_ != 0ULL; } - [[nodiscard]] std::uint64_t get_flushed_position() const noexcept; - [[nodiscard]] std::uint64_t get_ready_to_flush_position() const noexcept { + [[nodiscard]] std::uint64_t get_flushed_position() const; + [[nodiscard]] std::uint64_t get_ready_to_flush_position() const { return get_flushed_position() + last_transaction_boundary_position_in_event_buffer_; } diff --git a/src/binsrv/storage_core.cpp b/src/binsrv/storage_core.cpp index 33edd8e..cdca321 100644 --- a/src/binsrv/storage_core.cpp +++ b/src/binsrv/storage_core.cpp @@ -22,7 +22,9 @@ #include #include #include +#include #include +#include #include #include #include @@ -223,7 +225,9 @@ void storage_core::set_purged_gtids(const gtids::gtid_set &purged_gtids) { util::exception_location().raise( "cannot set purged GTIDs in position-based replication mode"); } - if (!is_empty()) { + + const std::unique_lock lock{mutex_}; + if (!is_empty_unsafe()) { util::exception_location().raise( "cannot set purged GTIDs in a non-empty storage"); } @@ -231,14 +235,20 @@ 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 + // method on an instance of basic_storage_backend that reads only data + // that was set only once during construction return backend_->get_description(); } [[nodiscard]] bool storage_core::is_in_gtid_replication_mode() const noexcept { + // no need to acquire the mutex as replication_mode_ is immutable after + // construction return replication_mode_ == replication_mode_type::gtid; } -[[nodiscard]] bool storage_core::is_binlog_open() const noexcept { +[[nodiscard]] bool storage_core::is_binlog_open() const { + const std::shared_lock lock{mutex_}; return backend_->is_stream_open(); } @@ -246,6 +256,8 @@ void storage_core::set_purged_gtids(const gtids::gtid_set &purged_gtids) { storage_core::open_binlog(const events::composite_binlog_name &binlog_name) { ensure_streaming_mode(); + const std::unique_lock lock{mutex_}; + auto result{open_binlog_status::opened_with_data_present}; // here we either create a new binlog file if its name is not presentin the @@ -258,12 +270,12 @@ storage_core::open_binlog(const events::composite_binlog_name &binlog_name) { // in the case when binlog exists, the name must be equal to the last item in // "binlog_records_" list and "position_" must be set to a non-zero value if (binlog_exists) { - if (binlog_name != get_current_binlog_name()) { + if (binlog_name != get_current_binlog_name_unsafe()) { util::exception_location().raise( "cannot open an existing binlog that is not the latest one for " "append"); } - if (get_flushed_position() == 0ULL) { + if (get_flushed_position_unsafe() == 0ULL) { util::exception_location().raise( "invalid position set when opening an existing binlog"); } @@ -288,24 +300,30 @@ void storage_core::write_event_block( events::seq_no_t block_max_sequence_number) { ensure_streaming_mode(); - write_data_to_stream(event_block_data, get_current_binlog_record().encryption, - get_current_binlog_record().size); - get_current_binlog_record().size += std::size(event_block_data); + const std::unique_lock lock{mutex_}; + + write_data_to_stream(event_block_data, + get_current_binlog_record_unsafe().encryption, + get_current_binlog_record_unsafe().size); + get_current_binlog_record_unsafe().size += std::size(event_block_data); if (is_in_gtid_replication_mode()) { - auto &optional_added_gtids{get_current_binlog_record().added_gtids}; + auto &optional_added_gtids{get_current_binlog_record_unsafe().added_gtids}; if (optional_added_gtids.has_value()) { *optional_added_gtids += block_gtids; } } - get_current_binlog_record().timestamps.add_range(block_timestamps); - get_current_binlog_record().last_sequence_number = block_max_sequence_number; + get_current_binlog_record_unsafe().timestamps.add_range(block_timestamps); + get_current_binlog_record_unsafe().last_sequence_number = + block_max_sequence_number; - save_binlog_metadata(get_current_binlog_record()); + save_binlog_metadata(get_current_binlog_record_unsafe()); } void storage_core::close_binlog() { ensure_streaming_mode(); + const std::unique_lock lock{mutex_}; + backend_->close_stream(); } @@ -313,7 +331,8 @@ void storage_core::close_binlog() { storage_core::purge_binlogs(const events::composite_binlog_name &target) { ensure_purging_mode(); - if (is_empty()) { + const std::unique_lock lock{mutex_}; + if (is_empty_unsafe()) { util::exception_location().raise( "cannot purge: binlog storage is empty"); } @@ -406,21 +425,30 @@ 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 + // 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 + // 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 + // 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 + // set only once during construction return encryption_format_.has_value() ? std::string{to_string_view(*encryption_format_)} : std::string{"encryption format is not set"}; @@ -520,7 +548,7 @@ void storage_core::ensure_purging_mode() const { gtids::optional_gtid_set previous_binlog_gtids{}; gtids::optional_gtid_set added_binlog_gtids{}; if (is_in_gtid_replication_mode()) { - previous_binlog_gtids = get_gtids(); + previous_binlog_gtids = get_gtids_unsafe(); added_binlog_gtids = gtids::gtid_set{}; } @@ -529,14 +557,14 @@ void storage_core::ensure_purging_mode() const { std::move(previous_binlog_gtids), std::move(added_binlog_gtids), util::ctime_timestamp_range{}, events::seq_no_t{}, std::move(encryption_record)); - save_binlog_metadata(get_current_binlog_record()); + save_binlog_metadata(get_current_binlog_record_unsafe()); save_binlog_index(); return open_binlog_status::created; } [[nodiscard]] open_binlog_status storage_core::open_existing_binlog_file_internal( std::uint64_t open_stream_offset) { - assert(get_flushed_position() == open_stream_offset); + assert(get_flushed_position_unsafe() == open_stream_offset); if (open_stream_offset >= events::magic_binlog_offset) { return open_stream_offset == events::magic_binlog_offset ? open_binlog_status::opened_at_magic_payload_offset @@ -545,8 +573,8 @@ storage_core::open_existing_binlog_file_internal( assert(open_stream_offset == 0ULL); write_data_to_stream(events::magic_binlog_payload, - get_current_binlog_record().encryption, 0ULL); - get_current_binlog_record().size = events::magic_binlog_offset; + get_current_binlog_record_unsafe().encryption, 0ULL); + get_current_binlog_record_unsafe().size = events::magic_binlog_offset; return open_binlog_status::opened_empty; } diff --git a/src/binsrv/storage_core.hpp b/src/binsrv/storage_core.hpp index 5a31bfb..f0a9105 100644 --- a/src/binsrv/storage_core.hpp +++ b/src/binsrv/storage_core.hpp @@ -19,6 +19,8 @@ #include "binsrv/storage_core_fwd.hpp" // IWYU pragma: export #include +#include +#include #include #include #include @@ -100,10 +102,14 @@ class [[nodiscard]] storage_core { [[nodiscard]] storage_construction_mode_type get_construction_mode() const noexcept { + // no mutex protection needed as this this method reads data + // set only once during construction return construction_mode_; } - [[nodiscard]] const gtids::gtid_set &get_purged_gtids() const noexcept { + // returning by value for thread-safety + [[nodiscard]] gtids::gtid_set get_purged_gtids() const { + const std::shared_lock lock{mutex_}; return purged_gtids_; } void set_purged_gtids(const gtids::gtid_set &purged_gtids); @@ -117,46 +123,36 @@ class [[nodiscard]] storage_core { } [[nodiscard]] bool is_in_gtid_replication_mode() const noexcept; - [[nodiscard]] const binlog_record_container & - get_binlog_records() const noexcept { + // returning by value for thread-safety + [[nodiscard]] binlog_record_container get_binlog_records() const { + const std::shared_lock lock{mutex_}; return binlog_records_; } - [[nodiscard]] bool is_empty() const noexcept { - return binlog_records_.empty(); + [[nodiscard]] bool is_empty() const { + const std::shared_lock lock{mutex_}; + return is_empty_unsafe(); } [[nodiscard]] events::composite_binlog_name get_current_binlog_name() const { - return is_empty() ? events::composite_binlog_name{} - : get_current_binlog_record().name; + const std::shared_lock lock{mutex_}; + return get_current_binlog_name_unsafe(); } [[nodiscard]] gtids::gtid_set get_gtids() const { - if (!is_in_gtid_replication_mode()) { - return {}; - } - - if (is_empty()) { - return get_purged_gtids(); - } - gtids::gtid_set result{}; - const auto &optional_previous_gtids{ - get_current_binlog_record().previous_gtids}; - if (optional_previous_gtids.has_value()) { - result = *optional_previous_gtids; - } - const auto &optional_added_gtids{get_current_binlog_record().added_gtids}; - if (optional_added_gtids.has_value()) { - result.add(*optional_added_gtids); - } - return result; + const std::shared_lock lock{mutex_}; + return get_gtids_unsafe(); } - [[nodiscard]] events::seq_no_t last_sequence_number() const noexcept { - return is_empty() ? 0ULL : get_current_binlog_record().last_sequence_number; + [[nodiscard]] events::seq_no_t get_last_sequence_number() const { + const std::shared_lock lock{mutex_}; + return is_empty_unsafe() + ? 0ULL + : get_current_binlog_record_unsafe().last_sequence_number; } - [[nodiscard]] std::uint64_t get_flushed_position() const noexcept { - return is_empty() ? 0ULL : get_current_binlog_record().size; + [[nodiscard]] std::uint64_t get_flushed_position() const { + const std::shared_lock lock{mutex_}; + return get_flushed_position_unsafe(); } - [[nodiscard]] bool is_binlog_open() const noexcept; + [[nodiscard]] bool is_binlog_open() const; [[nodiscard]] open_binlog_status open_binlog(const events::composite_binlog_name &binlog_name); @@ -187,6 +183,8 @@ class [[nodiscard]] storage_core { get_binlog_uri(const events::composite_binlog_name &binlog_name) const; [[nodiscard]] bool is_keyring_initialized() const noexcept { + // no mutex protection needed as this this method reads data + // set only once during construction return static_cast(keyring_); } [[nodiscard]] std::string get_keyring_description() const; @@ -194,10 +192,14 @@ class [[nodiscard]] storage_core { [[nodiscard]] std::string get_encryption_format_description() const; [[nodiscard]] bool has_active_kek() const noexcept { + // no mutex protection needed as this this method reads data + // set only once during construction return !active_kek_id_.empty(); } private: + mutable std::shared_mutex mutex_; + basic_logger_ptr logger_; storage_construction_mode_type construction_mode_; basic_keyring_ptr keyring_; @@ -210,6 +212,10 @@ class [[nodiscard]] storage_core { gtids::gtid_set purged_gtids_{}; binlog_record_container binlog_records_{}; + [[nodiscard]] bool is_empty_unsafe() const noexcept { + return binlog_records_.empty(); + } + void remove_temporary_objects(storage_object_name_container &object_names); void initialize_storage_encryption( @@ -219,12 +225,42 @@ class [[nodiscard]] storage_core { void ensure_purging_mode() const; [[nodiscard]] const binlog_record & - get_current_binlog_record() const noexcept { + get_current_binlog_record_unsafe() const noexcept { return binlog_records_.back(); } - [[nodiscard]] binlog_record &get_current_binlog_record() noexcept { + [[nodiscard]] binlog_record &get_current_binlog_record_unsafe() noexcept { return binlog_records_.back(); } + [[nodiscard]] events::composite_binlog_name + get_current_binlog_name_unsafe() const { + return is_empty_unsafe() ? events::composite_binlog_name{} + : get_current_binlog_record_unsafe().name; + } + [[nodiscard]] gtids::gtid_set get_gtids_unsafe() const { + if (!is_in_gtid_replication_mode()) { + return {}; + } + + if (is_empty_unsafe()) { + return purged_gtids_; + } + gtids::gtid_set result{}; + const auto &optional_previous_gtids{ + get_current_binlog_record_unsafe().previous_gtids}; + if (optional_previous_gtids.has_value()) { + result = *optional_previous_gtids; + } + const auto &optional_added_gtids{ + get_current_binlog_record_unsafe().added_gtids}; + if (optional_added_gtids.has_value()) { + result.add(*optional_added_gtids); + } + return result; + } + + [[nodiscard]] std::uint64_t get_flushed_position_unsafe() const noexcept { + return is_empty_unsafe() ? 0ULL : get_current_binlog_record_unsafe().size; + } [[nodiscard]] open_binlog_status open_new_binlog_file_internal( const events::composite_binlog_name &binlog_name);