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
16 changes: 6 additions & 10 deletions src/binsrv/storage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}

Expand All @@ -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();
Expand All @@ -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();
}

Expand All @@ -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 {
Expand Down Expand Up @@ -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();
}

Expand Down
15 changes: 7 additions & 8 deletions src/binsrv/storage.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -63,20 +63,19 @@ 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;

[[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_);
}

Expand All @@ -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);
Expand Down Expand Up @@ -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_;
}
Expand Down
62 changes: 45 additions & 17 deletions src/binsrv/storage_core.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,9 @@
#include <exception>
#include <filesystem>
#include <iterator>
#include <mutex>
#include <optional>
#include <shared_mutex>
#include <sstream>
#include <stdexcept>
#include <string>
Expand Down Expand Up @@ -223,29 +225,39 @@ void storage_core::set_purged_gtids(const gtids::gtid_set &purged_gtids) {
util::exception_location().raise<std::logic_error>(
"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<std::logic_error>(
"cannot set purged GTIDs in a non-empty storage");
}
purged_gtids_ = 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();
}

[[nodiscard]] open_binlog_status
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
Expand All @@ -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<std::logic_error>(
"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<std::logic_error>(
"invalid position set when opening an existing binlog");
}
Expand All @@ -288,32 +300,39 @@ 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();
}

[[nodiscard]] std::pair<binlog_record_container, std::string>
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<std::runtime_error>(
"cannot purge: binlog storage is empty");
}
Expand Down Expand Up @@ -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"};
Expand Down Expand Up @@ -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{};
}

Expand All @@ -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
Expand All @@ -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;
}

Expand Down
Loading
Loading