From 9ec512e504e661f666eeb892e3f4392db6b3397e Mon Sep 17 00:00:00 2001 From: Yura Sorokin Date: Wed, 23 Sep 2026 00:20:31 +0200 Subject: [PATCH] Rework storage layer abstractions to support simultaneous reads and writes (part 2) https://perconadev.atlassian.net/browse/PBS-6 This is a pre-requisite commit with mostly refactorings: * Extracted "unbuffered" part of the 'binsrv::storage' class (that includes binlog_record_container and basic_storage_backend_ptr) into 'binsrv::storage_core'. * 'binsrv::storage' class interface reimplemented with this new 'binsrv::storage_core'.. * 'binsrv::storage::binlog_encryption_record' class-level struct moved to namespace level 'binsrv::binlog_encryption_record'. * 'binsrv::storage::binlog_record' class-level struct moved to namespace level 'binsrv::binlog_record'. The idea behind this refactoring is to separate methods that need to be protected by a mutex from those that do not. --- CMakeLists.txt | 4 + src/binsrv/storage.cpp | 961 ++--------------- src/binsrv/storage.hpp | 158 +-- src/binsrv/storage_core.cpp | 987 ++++++++++++++++++ src/binsrv/storage_core.hpp | 267 +++++ src/binsrv/storage_core_fwd.hpp | 51 + src/binsrv/storage_fwd.hpp | 14 - src/operations/collector_context.cpp | 1 + src/operations/fetch_operation.cpp | 1 + src/operations/list_operation.cpp | 1 + src/operations/model_helpers.cpp | 10 +- src/operations/model_helpers.hpp | 6 +- src/operations/pull_operation.cpp | 1 + src/operations/purge_binlogs_operation.cpp | 1 + .../search_by_gtid_set_operation.cpp | 1 + .../search_by_timestamp_operation.cpp | 1 + 16 files changed, 1397 insertions(+), 1068 deletions(-) create mode 100644 src/binsrv/storage_core.cpp create mode 100644 src/binsrv/storage_core.hpp create mode 100644 src/binsrv/storage_core_fwd.hpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 65f2555..8ee3463 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -620,6 +620,10 @@ set(binsrv_source_files src/binsrv/storage.hpp src/binsrv/storage.cpp + src/binsrv/storage_core_fwd.hpp + src/binsrv/storage_core.hpp + src/binsrv/storage_core.cpp + src/binsrv/storage_backend_factory.hpp src/binsrv/storage_backend_factory.cpp diff --git a/src/binsrv/storage.cpp b/src/binsrv/storage.cpp index 4640efe..634bb3d 100644 --- a/src/binsrv/storage.cpp +++ b/src/binsrv/storage.cpp @@ -15,225 +15,43 @@ #include "binsrv/storage.hpp" -#include #include #include #include #include -#include -#include #include -#include -#include +#include #include #include #include #include #include -#include "binsrv/basic_keyring.hpp" #include "binsrv/basic_logger.hpp" -#include "binsrv/basic_storage_backend.hpp" -#include "binsrv/binlog_file_metadata.hpp" -#include "binsrv/encryption_format_type.hpp" -#include "binsrv/keyring_factory.hpp" -#include "binsrv/keyring_record.hpp" -#include "binsrv/log_severity.hpp" #include "binsrv/main_config.hpp" -#include "binsrv/replication_config.hpp" -#include "binsrv/replication_mode_type.hpp" -#include "binsrv/storage_backend_factory.hpp" -#include "binsrv/storage_config.hpp" -#include "binsrv/storage_metadata.hpp" +#include "binsrv/replication_mode_type_fwd.hpp" +#include "binsrv/storage_core.hpp" #include "binsrv/events/common_types.hpp" #include "binsrv/events/composite_binlog_name.hpp" -#include "binsrv/events/protocol_traits_fwd.hpp" #include "binsrv/gtids/gtid.hpp" #include "binsrv/gtids/gtid_set.hpp" -#include "binsrv/models/binlog_file_encryption_record.hpp" - -#include "opensslpp/cipher_context.hpp" -#include "opensslpp/crypto_rng.hpp" - #include "util/byte_span.hpp" #include "util/ctime_timestamp.hpp" #include "util/exception_location_helpers.hpp" namespace binsrv { -[[nodiscard]] models::binlog_file_encryption_record -storage::binlog_encryption_record::to_model( - const storage::binlog_encryption_record &record) { - models::binlog_file_encryption_record model{}; - - auto &file_key_envelope{model.get<"file_key_envelope">()}; - file_key_envelope.get<"kek_id">() = record.kek_id; - file_key_envelope.get<"data_hex">() = record.file_key_encrypted_with_kek; - if (record.iv_for_file_key_encryption.has_value()) { - file_key_envelope.get<"iv_hex">() = *record.iv_for_file_key_encryption; - } - if (record.tag_of_file_key_encryption.has_value()) { - file_key_envelope.get<"tag_hex">() = *record.tag_of_file_key_encryption; - } - - auto &file_data_envelope{model.get<"file_data_envelope">()}; - file_data_envelope.get<"cipher">() = record.data_cipher; - file_data_envelope.get<"iv_hex">() = record.iv_for_data_encryption; - if (record.tag_of_data_encryption.has_value()) { - file_data_envelope.get<"tag_hex">() = *record.tag_of_data_encryption; - } - return model; -} - -[[nodiscard]] storage::binlog_encryption_record -storage::binlog_encryption_record::from_model( - const models::binlog_file_encryption_record &model) { - binlog_encryption_record record{}; - - const auto &file_key_envelope{model.get<"file_key_envelope">()}; - record.kek_id = file_key_envelope.get<"kek_id">(); - const auto file_key_raw{file_key_envelope.get<"data_hex">().get_data()}; - record.file_key_encrypted_with_kek.assign(std::cbegin(file_key_raw), - std::cend(file_key_raw)); - if (file_key_envelope.get<"iv_hex">().has_value()) { - const auto file_key_iv_raw{file_key_envelope.get<"iv_hex">()->get_data()}; - record.iv_for_file_key_encryption.emplace(std::cbegin(file_key_iv_raw), - std::cend(file_key_iv_raw)); - } - if (file_key_envelope.get<"tag_hex">().has_value()) { - const auto file_key_tag_raw{file_key_envelope.get<"tag_hex">()->get_data()}; - record.tag_of_file_key_encryption.emplace(std::cbegin(file_key_tag_raw), - std::cend(file_key_tag_raw)); - } - - const auto &file_data_envelope{model.get<"file_data_envelope">()}; - record.data_cipher = file_data_envelope.get<"cipher">(); - const auto file_data_iv_raw{file_data_envelope.get<"iv_hex">().get_data()}; - record.iv_for_data_encryption.assign(std::cbegin(file_data_iv_raw), - std::cend(file_data_iv_raw)); - if (file_data_envelope.get<"tag_hex">().has_value()) { - const auto file_data_tag_raw{ - file_data_envelope.get<"tag_hex">()->get_data()}; - record.tag_of_data_encryption.emplace(std::cbegin(file_data_tag_raw), - std::cend(file_data_tag_raw)); - } - return record; -} - storage::storage(basic_logger_ptr logger, const main_config &config, storage_construction_mode_type construction_mode) - : logger_{std::move(logger)}, construction_mode_{construction_mode}, - backend_{} { - assert(logger_); - - const auto &replication_config{config.root().get<"replication">()}; - // we need a copy of replication mode as replication_mode_ will be - // overwritten by load_metadata() later - const auto replication_mode{replication_config.get<"mode">()}; - replication_mode_ = replication_mode; - - const auto &storage_config{config.root().get<"storage">()}; - - const auto &checkpoint_size_opt{storage_config.get<"checkpoint_size">()}; - if (checkpoint_size_opt.has_value()) { - checkpoint_size_bytes_ = checkpoint_size_opt->get_value(); - } - - const auto &checkpoint_interval_opt{ - storage_config.get<"checkpoint_interval">()}; - if (checkpoint_interval_opt.has_value()) { - checkpoint_interval_seconds_ = - std::chrono::seconds{checkpoint_interval_opt->get_value()}; - } - - const auto &keyring_config{config.root().get<"keyring">()}; - if (keyring_config.has_value()) { - keyring_ = keyring_factory::create(keyring_config->get<"uri">()); - } - const auto &encryption_config{storage_config.get<"encryption">()}; - initialize_storage_encryption(encryption_config); - - backend_ = storage_backend_factory::create(storage_config); - - auto storage_objects{backend_->list_objects()}; - remove_temporary_objects(storage_objects); - - if (storage_objects.empty()) { - // initialized on a new / empty storage - just save metadata and return - if (construction_mode_ == storage_construction_mode_type::streaming) { - save_metadata(); - } - return; - } - - const auto metadata_it{std::as_const(storage_objects).find(metadata_name)}; - if (metadata_it == std::cend(storage_objects)) { - util::exception_location().raise( - "storage is not empty but does not contain metadata"); - } - storage_objects.erase(metadata_it); - - // as load_metadata() will be updating 'encryption_format_', saving it here - // to use for validation later - const auto encryption_format{encryption_format_}; - load_metadata(); - validate_metadata(replication_mode, encryption_format); - // in case when storage metadata file is present and did not have encryption - // format specified, but it is set in the configuration file, we need to - // update storage metadata - if (!encryption_format_.has_value() && encryption_format.has_value()) { - encryption_format_ = encryption_format; - save_metadata(); - } - - // if after metadata erasure 'storage_objects' is empty, then this means - // that it has only metadata in it that passes validation and we can - // consider it as an initialized empty storage, so just return - if (storage_objects.empty()) { - return; - } - - const auto binlog_index_it{storage_objects.find(default_binlog_index_name)}; - if (binlog_index_it == std::cend(storage_objects)) { - // as binlog index file is created after the very first binlog data file - // and its metadata file are created, the concurrent query-only operation - // may intrude exactly between these two steps, so for query-only operations - // we should not consider the absence of binlog index file as an error - if (construction_mode == storage_construction_mode_type::querying_only) { - return; - } - util::exception_location().raise( - "storage is not empty but does not contain binlog index"); - } - storage_objects.erase(binlog_index_it); - - // extracting all binlog file metadata files into a separate container - storage_object_name_container storage_metadata_objects; - for (auto storage_object_it{std::cbegin(storage_objects)}; - storage_object_it != std::cend(storage_objects);) { - const std::filesystem::path object_name{storage_object_it->first}; - if (object_name.has_extension() && - object_name.extension() == binlog_metadata_extension) { - auto object_node = storage_objects.extract(storage_object_it++); - storage_metadata_objects.insert(std::move(object_node)); - } else { - ++storage_object_it; - } - } - load_binlog_index(); - validate_binlog_index(storage_objects); - - load_and_validate_binlog_metadata_set(storage_objects, - storage_metadata_objects); - assert(!binlog_records_.front().added_gtids.has_value() || - purged_gtids_ == binlog_records_.front().added_gtids); -} + : core_{std::make_unique(std::move(logger), config, + construction_mode)} {} storage::~storage() { - if (construction_mode_ == storage_construction_mode_type::streaming) { + if (core_->get_construction_mode() == + storage_construction_mode_type::streaming) { // bugprone-empty-catch should not be that strict in destructors try { flush_event_buffer(); @@ -242,69 +60,59 @@ storage::~storage() { } } +[[nodiscard]] const gtids::gtid_set & +storage::get_purged_gtids() const noexcept { + return core_->get_purged_gtids(); +} + void storage::set_purged_gtids(const gtids::gtid_set &purged_gtids) { - if (!is_in_gtid_replication_mode()) { - util::exception_location().raise( - "cannot set purged GTIDs in position-based replication mode"); - } - if (!is_empty()) { - util::exception_location().raise( - "cannot set purged GTIDs in a non-empty storage"); - } - purged_gtids_ = purged_gtids; + core_->set_purged_gtids(purged_gtids); } [[nodiscard]] std::string storage::get_backend_description() const { - return backend_->get_description(); + return core_->get_backend_description(); +} + +[[nodiscard]] replication_mode_type +storage::get_replication_mode() const noexcept { + return core_->get_replication_mode(); } [[nodiscard]] bool storage::is_in_gtid_replication_mode() const noexcept { - return replication_mode_ == replication_mode_type::gtid; + return core_->is_in_gtid_replication_mode(); } -[[nodiscard]] bool storage::is_binlog_open() const noexcept { - return backend_->is_stream_open(); +[[nodiscard]] const binlog_record_container & +storage::get_binlog_records() const noexcept { + return core_->get_binlog_records(); } -[[nodiscard]] open_binlog_status -storage::open_binlog(const events::composite_binlog_name &binlog_name) { - ensure_streaming_mode(); +[[nodiscard]] bool storage::is_empty() const noexcept { + return core_->is_empty(); +} - auto result{open_binlog_status::opened_with_data_present}; +[[nodiscard]] gtids::gtid_set storage::get_gtids() const { + return core_->get_gtids(); +} - // here we either create a new binlog file if its name is not presentin the - // "binlog_records_", or we open an existing one and append to it, in which - // case we need to make sure that the current position is properly set - const bool binlog_exists{ - std::ranges::find(std::as_const(binlog_records_), binlog_name, - &binlog_record::name) != std::cend(binlog_records_)}; +[[nodiscard]] events::composite_binlog_name +storage::get_current_binlog_name() const { + return core_->get_current_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()) { - util::exception_location().raise( - "cannot open an existing binlog that is not the latest one for " - "append"); - } - if (get_current_position() == 0ULL) { - util::exception_location().raise( - "invalid position set when opening an existing binlog"); - } - } +[[nodiscard]] bool storage::is_binlog_open() const noexcept { + return core_->is_binlog_open(); +} - const auto mode{binlog_exists ? storage_backend_open_stream_mode::append - : storage_backend_open_stream_mode::create}; - const auto open_stream_offset{backend_->open_stream(binlog_name.str(), mode)}; +[[nodiscard]] open_binlog_status +storage::open_binlog(const events::composite_binlog_name &binlog_name) { + auto result{core_->open_binlog(binlog_name)}; - if (binlog_exists) { - result = open_existing_binlog_file_internal(open_stream_offset); - ready_to_flush_last_sequence_number_ = - binlog_records_.back().last_sequence_number; + if (result != open_binlog_status::created) { + ready_to_flush_last_sequence_number_ = core_->last_sequence_number(); incomplete_transaction_last_sequence_number_ = ready_to_flush_last_sequence_number_; } else { - result = open_new_binlog_file_internal(binlog_name); ready_to_flush_last_sequence_number_ = 0ULL; incomplete_transaction_last_sequence_number_ = 0ULL; } @@ -330,6 +138,9 @@ void storage::write_event(util::const_byte_span event_data, event_buffer_.insert(std::end(event_buffer_), std::cbegin(event_data), std::cend(event_data)); incomplete_transaction_timestamps_.add_timestamp(event_timestamp); + // 0 has a special meaning here - it indicates that current event is neither + // GTID_LOG, nor ANONYMOUS_GTID_LOG, nor GTID_TAGGED_LOG and does not have + // last sequence number associated with it. if (transaction_sequence_number != 0ULL) { incomplete_transaction_last_sequence_number_ = transaction_sequence_number; } @@ -387,7 +198,7 @@ void storage::close_binlog() { event_buffer_.clear(); event_buffer_.shrink_to_fit(); - backend_->close_stream(); + core_->close_binlog(); update_last_checkpoint_info(); } @@ -408,203 +219,26 @@ void storage::flush_event_buffer() { } } -[[nodiscard]] std::pair +[[nodiscard]] std::pair storage::purge_binlogs(const events::composite_binlog_name &target) { - ensure_purging_mode(); - - if (is_empty()) { - util::exception_location().raise( - "cannot purge: binlog storage is empty"); - } - const auto &front_base_name{binlog_records_.front().name.get_base_name()}; - if (target.get_base_name() != front_base_name) { - util::exception_location().raise( - "cannot purge: target binlog name has a different base name than " - "the binlog records in the storage"); - } - const auto target_it{std::ranges::find(std::as_const(binlog_records_), target, - &binlog_record::name)}; - if (target_it == std::cend(binlog_records_)) { - util::exception_location().raise( - "cannot purge: target binlog name is not present in the storage"); - } - // refuse to purge the current tail: emptying the storage would lose - // the resume position (current binlog name / position in position - // mode, executed GTID set in GTID mode) and force the next 'fetch' / - // 'pull' to re-stream from the very beginning of the source's - // retained binlog history. - if (target_it == std::prev(std::cend(binlog_records_))) { - util::exception_location().raise( - "cannot purge: target is the current tail binlog file; at least " - "one binlog file must remain in the storage to preserve the " - "resume position"); - } - - // step 1: extract the prefix [begin, target_it + 1) - this - // becomes the set of records we are going to drop on disk; the - // returned vector preserves the original order so the caller can - // use it directly to produce a response - const auto victim_count{static_cast( - std::distance(std::cbegin(binlog_records_), target_it) + 1)}; - binlog_record_container removed_records; - removed_records.reserve(victim_count); - std::move(std::begin(binlog_records_), - std::begin(binlog_records_) + - static_cast(victim_count), - std::back_inserter(removed_records)); - binlog_records_.erase(std::begin(binlog_records_), - std::begin(binlog_records_) + - static_cast(victim_count)); - - // step 2: rewrite the binlog index from the surviving records left - // in 'binlog_records_' after step 1 (always non-empty thanks to the - // tail-refusal guard above). 'save_binlog_index' goes through the - // backend's atomic-overwrite 'put_object', so from this point on - // the purge is considered committed - any subsequent failure - // leaves the storage in an inconsistent state (leftover payload / - // metadata files no longer referenced by the index) that the - // constructor's existing validators will refuse to open on next - // startup. - save_binlog_index(); - - // step 3: best-effort removal of the victim payload + metadata - // objects; any failure here is intentionally swallowed - the index - // has already been committed and reporting a "file could not be - // removed" error to the caller would falsely suggest that the - // purge itself failed; the resulting leftovers will trip the - // constructor's validators on next startup. - // We materialise the (metadata + payload) names for every victim - // into a single batch and hand it to 'basic_storage_backend:: - // remove_objects', which runs the backend's durability barrier - // exactly once at the end of the batch - so the whole batch - // amortises to a single fsync(2) on the local filesystem backend - // (and a no-op on S3) instead of O(N) syncs. - std::vector victim_object_names; - victim_object_names.reserve(std::size(removed_records) * 2U); - for (const auto &victim : removed_records) { - victim_object_names.emplace_back( - generate_binlog_metadata_name(victim.name)); - victim_object_names.emplace_back(victim.name.str()); - } - std::string cleanup_warning_message; - try { - backend_->remove_objects(victim_object_names); - } catch (const std::exception &e) { - // 'remove_objects' re-raises the first per-name failure (if any) - // after running the durability barrier; we do not propagate it - // to the caller because the index has already been committed - // and any leftover payload/metadata files will be picked up by - // the constructor's validators on the next startup. We just - // capture the underlying message so the caller can surface it - // under a 'warning' status in the JSON response. - cleanup_warning_message = e.what(); - } - - return {std::move(removed_records), std::move(cleanup_warning_message)}; + return core_->purge_binlogs(target); } [[nodiscard]] std::string storage::get_binlog_uri( const events::composite_binlog_name &binlog_name) const { - return backend_->get_object_uri(binlog_name.str()); + return core_->get_binlog_uri(binlog_name); } [[nodiscard]] std::string storage::get_keyring_description() const { - return is_keyring_initialized() ? keyring_->get_description() - : "keyring is not initialized"; + return core_->get_keyring_description(); } [[nodiscard]] std::string storage::get_active_kek_description() const { - return has_active_kek() ? keyring_->get_key(active_kek_id_).get_description() - : "active KEK is not set"; + return core_->get_active_kek_description(); } [[nodiscard]] std::string storage::get_encryption_format_description() const { - return encryption_format_.has_value() - ? std::string{to_string_view(*encryption_format_)} - : std::string{"encryption format is not set"}; -} - -void storage::remove_temporary_objects( - storage_object_name_container &object_names) { - using remove_object_container = std::vector; - remove_object_container remove_objects; - for (auto it{std::begin(object_names)}; it != std::end(object_names);) { - const std::filesystem::path object_name{it->first}; - if (object_name.has_extension() && - object_name.extension() == tmp_storage_object_suffix) { - auto object_node = object_names.extract(it++); - remove_objects.emplace_back(std::move(object_node.key())); - logger_->log_format(log_severity::warning, - "found temporary storage object '{}' left after " - "improper shutdown", - object_name.string()); - } else { - ++it; - } - } - if (remove_objects.empty()) { - return; - } - - // for querying-only mode we do not perform any modifying operations - just - // clear the 'object_names' container and return - if (construction_mode_ == storage_construction_mode_type::querying_only) { - return; - } - - backend_->remove_objects(remove_objects); - logger_->log_format(log_severity::warning, - "removed {} temporary storage object(s) left after " - "improper shutdown", - remove_objects.size()); -} - -void storage::initialize_storage_encryption( - const optional_encryption_config &encryption_config) { - if (encryption_config.has_value()) { - encryption_format_ = encryption_config->get<"format">(); - if (!is_keyring_initialized()) { - util::exception_location().raise( - "encryption is enabled but keyring is not initialized"); - } - const auto &kek_id{encryption_config->get<"kek_id">()}; - if (!keyring_->contains(kek_id)) { - util::exception_location().raise( - "keyring does not contain the specified KEK ID"); - } - active_kek_id_ = kek_id; - active_data_cipher_ = encryption_config->get<"cipher">(); - - // make sure that random file keys (of length that corresponds to the - // active data cipher) can be encrypted with the active KEK - - // for instance, if active data cipher is AES-192-CRT (key length 24 - // bytes), then the active KEK cannot be of ECB or CBC mode as these - // ciphers can only encrypt data of length that is a multiple of the - // block size (16 bytes) - const auto &keyring_record{keyring_->get_key(active_kek_id_)}; - if (opensslpp::cipher_context::get_key_size_in_bytes(active_data_cipher_) % - opensslpp::cipher_context::get_block_size_in_bytes( - keyring_record.get<"cipher">()) != - 0U) { - util::exception_location().raise( - "active data cipher key length is not compatible with the active " - "KEK cipher block size"); - } - } -} - -void storage::ensure_streaming_mode() const { - if (construction_mode_ != storage_construction_mode_type::streaming) { - util::exception_location().raise( - "operation requires storage to be constructed in streaming mode"); - } -} - -void storage::ensure_purging_mode() const { - if (construction_mode_ != storage_construction_mode_type::purging) { - util::exception_location().raise( - "operation requires storage to be constructed in purging mode"); - } + return core_->get_encryption_format_description(); } void storage::update_last_checkpoint_info() { @@ -616,44 +250,8 @@ void storage::update_last_checkpoint_info() { } } -[[nodiscard]] open_binlog_status storage::open_new_binlog_file_internal( - const events::composite_binlog_name &binlog_name) { - - auto encryption_record{generate_binlog_encryption_record()}; - // writing the magic binlog footprint only if this is a newly - // created file - write_data_to_stream(events::magic_binlog_payload, encryption_record, 0ULL); - - 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(); - added_binlog_gtids = gtids::gtid_set{}; - } - - binlog_records_.emplace_back( - binlog_name, events::magic_binlog_offset, - 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_index(); - return open_binlog_status::created; -} -[[nodiscard]] open_binlog_status -storage::open_existing_binlog_file_internal(std::uint64_t open_stream_offset) { - assert(get_current_position() == 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 - : open_binlog_status::opened_with_data_present; - } - 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; - return open_binlog_status::opened_empty; +[[nodiscard]] std::uint64_t storage::get_flushed_position() const noexcept { + return core_->get_flushed_position(); } void storage::flush_event_buffer_internal() { @@ -666,22 +264,9 @@ void storage::flush_event_buffer_internal() { last_transaction_boundary_position_in_event_buffer_}; // writing bytes from // the beginning of the event buffer - write_data_to_stream(transactions_data, - get_current_binlog_record().encryption, - get_current_binlog_record().size); - get_current_binlog_record().size += - last_transaction_boundary_position_in_event_buffer_; - if (is_in_gtid_replication_mode()) { - auto &optional_added_gtids{get_current_binlog_record().added_gtids}; - if (optional_added_gtids.has_value()) { - *optional_added_gtids += gtids_in_event_buffer_; - } - } - get_current_binlog_record().timestamps.add_range(ready_to_flush_timestamps_); - get_current_binlog_record().last_sequence_number = - ready_to_flush_last_sequence_number_; - - save_binlog_metadata(get_current_binlog_record()); + core_->write_event_block(transactions_data, gtids_in_event_buffer_, + ready_to_flush_timestamps_, + ready_to_flush_last_sequence_number_); const auto begin_it{std::cbegin(event_buffer_)}; const auto portion_it{std::next( @@ -697,438 +282,12 @@ void storage::flush_event_buffer_internal() { ready_to_flush_timestamps_.clear(); } -void storage::load_binlog_index() { - const auto index_content{backend_->get_object(default_binlog_index_name)}; - // opening in text mode - std::istringstream index_iss{index_content}; - std::string current_line; - while (std::getline(index_iss, current_line)) { - if (current_line.empty()) { - continue; - } - const std::filesystem::path current_binlog_path{current_line}; - if (current_binlog_path.parent_path() != default_binlog_index_entry_path) { - util::exception_location().raise( - "binlog index contains an entry that has an invalid path"); - } - auto current_binlog_name{current_binlog_path.filename().string()}; - - if (current_binlog_name == default_binlog_index_name) { - util::exception_location().raise( - "binlog index contains a reference to the binlog index name"); - } - const auto current_binlog_name_parsed{ - events::composite_binlog_name::parse(current_binlog_name)}; - if (std::ranges::find(std::as_const(binlog_records_), - current_binlog_name_parsed, - &binlog_record::name) != std::cend(binlog_records_)) { - util::exception_location().raise( - "binlog index contains a duplicate entry"); - } - gtids::optional_gtid_set previous_binlog_gtids{}; - gtids::optional_gtid_set added_binlog_gtids{}; - if (is_in_gtid_replication_mode()) { - previous_binlog_gtids = gtids::gtid_set{}; - added_binlog_gtids = gtids::gtid_set{}; - } - binlog_records_.emplace_back( - current_binlog_name_parsed, 0ULL, std::move(previous_binlog_gtids), - std::move(added_binlog_gtids), util::ctime_timestamp_range{}); - } -} - -void storage::validate_binlog_index( - const storage_object_name_container &object_names) const { - // in the querying_only mode we allow discrepancies between the binlog index - // and the actual objects in the storage - if (construction_mode_ == storage_construction_mode_type::querying_only) { - return; - } - - for (auto const &record : binlog_records_) { - if (!object_names.contains(record.name.str())) { - util::exception_location().raise( - "binlog index contains a reference to a non-existing object"); - } - } - - if (std::size(object_names) != std::size(binlog_records_)) { - util::exception_location().raise( - "storage contains an object that is not " - "referenced in the binlog index"); - } - - // TODO: add integrity checks (parsing + checksumming) for the binlog - // files in the index -} - -void storage::save_binlog_index() const { - std::ostringstream oss; - for (const auto &record : binlog_records_) { - std::filesystem::path binlog_path{default_binlog_index_entry_path}; - binlog_path /= record.name.str(); - oss << binlog_path.generic_string() << '\n'; - } - const auto content{oss.str()}; - backend_->put_object(default_binlog_index_name, - util::as_const_byte_span(content)); -} - -void storage::load_metadata() { - const auto metadata_content{backend_->get_object(metadata_name)}; - const storage_metadata metadata{metadata_content}; - replication_mode_ = metadata.root().get<"mode">(); - encryption_format_ = metadata.root().get<"encryption">(); -} - -void storage::validate_metadata( - replication_mode_type replication_mode, - const optional_encryption_format_type &encryption_format) const { - if (replication_mode != replication_mode_) { - util::exception_location().raise( - "replication mode provided to initialize storage differs from the one " - "stored in metadata"); - } - - if (encryption_format_.has_value() && encryption_format.has_value() && - *encryption_format_ != *encryption_format) { - // if both encryption formats (the existing one loaded from the storage - // metadata file and a new one specified in the configuration file) are - // present and have different values, then this is an error +void storage::ensure_streaming_mode() const { + if (core_->get_construction_mode() != + storage_construction_mode_type::streaming) { util::exception_location().raise( - "storage encryption format provided to initialize storage differs from " - "the one stored in metadata"); - } -} - -void storage::save_metadata() const { - storage_metadata metadata{}; - metadata.root().get<"mode">() = replication_mode_; - metadata.root().get<"encryption">() = encryption_format_; - const auto content{metadata.str()}; - backend_->put_object(metadata_name, util::as_const_byte_span(content)); -} - -[[nodiscard]] std::string storage::generate_binlog_metadata_name( - const events::composite_binlog_name &binlog_name) { - std::string binlog_metadata_name{binlog_name.str()}; - binlog_metadata_name += storage::binlog_metadata_extension; - return binlog_metadata_name; -} - -[[nodiscard]] storage::binlog_record storage::load_binlog_metadata( - const events::composite_binlog_name &binlog_name) const { - const auto content{ - backend_->get_object(generate_binlog_metadata_name(binlog_name))}; - binlog_file_metadata metadata{content}; - - const auto &optional_encryption_metadata{metadata.root().get<"encryption">()}; - return binlog_record{ - .name = binlog_name, - .size = metadata.root().get<"size">(), - .previous_gtids = metadata.root().get<"previous_gtids">(), - .added_gtids = metadata.root().get<"added_gtids">(), - .timestamps = {metadata.root().get<"min_timestamp">(), - metadata.root().get<"max_timestamp">()}, - .last_sequence_number = metadata.root().get<"last_sequence_number">(), - .encryption = optional_encryption_metadata.has_value() - ? binlog_encryption_record::from_model( - *optional_encryption_metadata) - : optional_binlog_encryption_record{}}; -} - -void storage::validate_binlog_metadata(const binlog_record &record) const { - if (is_in_gtid_replication_mode()) { - if (!record.previous_gtids.has_value()) { - util::exception_location().raise( - "missing previous GTID set in the binlog metadata while in GTID " - "replication " - "mode"); - } - if (!record.added_gtids.has_value()) { - util::exception_location().raise( - "missing added GTID set in the binlog metadata while in GTID " - "replication " - "mode"); - } - } else { - if (record.previous_gtids.has_value()) { - util::exception_location().raise( - "found previous GTID set in the binlog metadata while in position " - "replication mode"); - } - if (record.added_gtids.has_value()) { - util::exception_location().raise( - "found added GTID set in the binlog metadata while in position " - "replication mode"); - } - } - // make sure that if the encryption record is present in the binlog - // metadata, keyring must be initialized and contain the KEK with the - // ID specified in the encryption record - if (record.encryption.has_value()) { - if (!is_keyring_initialized()) { - util::exception_location().raise( - "found encryption record in the binlog metadata but keyring is not " - "initialized"); - } - if (!keyring_->contains(record.encryption->kek_id)) { - util::exception_location().raise( - "found encryption record in the binlog metadata but keyring does not " - "contain the specified KEK ID"); - } - } -} - -void storage::save_binlog_metadata(const binlog_record &record) const { - binlog_file_metadata metadata{}; - metadata.root().get<"size">() = record.size; - metadata.root().get<"previous_gtids">() = record.previous_gtids; - metadata.root().get<"added_gtids">() = record.added_gtids; - metadata.root().get<"min_timestamp">() = - util::ctime_timestamp{record.timestamps.get_min_timestamp()}; - metadata.root().get<"max_timestamp">() = - util::ctime_timestamp{record.timestamps.get_max_timestamp()}; - metadata.root().get<"last_sequence_number">() = record.last_sequence_number; - const auto &record_encryption{record.encryption}; - if (record_encryption.has_value()) { - metadata.root().get<"encryption">() = - binlog_encryption_record::to_model(*record_encryption); - } - const auto content{metadata.str()}; - backend_->put_object(generate_binlog_metadata_name(record.name), - util::as_const_byte_span(content)); -} - -void storage::load_and_validate_binlog_metadata_set( - // NOLINTNEXTLINE(bugprone-easily-swappable-parameters) - const storage_object_name_container &object_names, - const storage_object_name_container &object_metadata_names) { - auto record_it{std::begin(binlog_records_)}; - while (record_it != std::end(binlog_records_)) { - storage::binlog_record loaded_binlog_metadata{}; - try { - const auto binlog_metadata_name{ - generate_binlog_metadata_name(record_it->name)}; - if (!object_metadata_names.contains(binlog_metadata_name)) { - util::exception_location().raise( - "missing metadata for a binlog listed in the binlog index"); - } - loaded_binlog_metadata = load_binlog_metadata(record_it->name); - } catch (const std::exception &) { - if (construction_mode_ == storage_construction_mode_type::querying_only) { - // in the querying_only mode we just skip invalid metadata and the - // corresponding binlog file - this allows to query an otherwise - // unusable storage and retrieve information about valid binlog files - // from it, which can be useful for debugging / forensics purposes - record_it = binlog_records_.erase(record_it); - continue; - } - throw; - } - validate_binlog_metadata(loaded_binlog_metadata); - // validating binlog size from the metadata only makes sense if we are not - // in the querying_only mode - if (construction_mode_ != storage_construction_mode_type::querying_only) { - const auto binlog_file_name{record_it->name.str()}; - const auto actual_binlog_file_size{object_names.at(binlog_file_name)}; - // validating that the size stored in the metadata matches the actual size - if (loaded_binlog_metadata.size != actual_binlog_file_size) { - // in case when Binlog Server process was not properly shut down - // there is a chance that there will be mismatch between the actual - // binlog data file size and the 'size' field in the binlog metadata - - // if this mismatch was found in the most recent binlog file, we can - // perform automatic recovery (truncating binlog data file content to - // the size from the metadata) - if (std::next(record_it) != std::end(binlog_records_)) { - util::exception_location().raise( - "size from the binlog metadata does not match the actual binlog " - "size"); - } - // performing recovery - if (loaded_binlog_metadata.size > actual_binlog_file_size) { - util::exception_location().raise( - "cannot perform recovery - size from the binlog metadata is " - "bigger than the actual binlog file size"); - } - backend_->resize_object(binlog_file_name, loaded_binlog_metadata.size); - logger_->log_format(log_severity::warning, - "recovered binlog file '{}' by truncating it to " - "the size from the metadata", - binlog_file_name); - } - } - *record_it = std::move(loaded_binlog_metadata); - ++record_it; - } - // after this loop position_ and gtids_ should store the values from the last - // binlog file metadata - - if (construction_mode_ != storage_construction_mode_type::querying_only) { - if (std::size(object_metadata_names) != std::size(binlog_records_)) { - util::exception_location().raise( - "found metadata for a non-existing binlog"); - } - } - - // if we are in GTID replication mode, then we can consider GTIDs from the - // first binlog metadata as purged GTIDs for the whole storage - const auto &optional_added_gtids{binlog_records_.front().added_gtids}; - if (optional_added_gtids.has_value()) { - purged_gtids_ = *optional_added_gtids; - } -} - -[[nodiscard]] storage::optional_binlog_encryption_record -storage::generate_binlog_encryption_record() const { - if (!has_active_kek()) { - return std::nullopt; - } - - // we identify the KEK record in the keyring by the active KEK ID, - // specified in the main configuration file - // ('' parameter) - const auto &keyring_record{keyring_->get_key(active_kek_id_)}; - - // identifying the the cipher name and the key data from the - // keyring record - this data will be used to encrypt random file - // keys generated for new binlog data files - const auto &kek_cipher{keyring_record.get<"cipher">()}; - const auto &kek{keyring_record.get<"data_hex">().get_data()}; - - // identify the size of the IV that will be used for file key - // encryption based on the KEK cipher; if the KEK cipher is in ECB mode, then - // the IV is not used and its size will be 0 - const auto iv_size_for_file_key_encryption{ - opensslpp::cipher_context::get_iv_size_in_bytes(kek_cipher)}; - util::optional_hex_value_storage iv_for_file_key_encryption{}; - util::const_byte_span iv_for_file_key_encryption_v{}; - if (iv_size_for_file_key_encryption != 0U) { - // generating random IV for file key encryption - iv_for_file_key_encryption.emplace(iv_size_for_file_key_encryption); - opensslpp::crypto_rng::generate(*iv_for_file_key_encryption); - iv_for_file_key_encryption_v = *iv_for_file_key_encryption; - } - - // identify the size of the file key based on the active data cipher - const auto file_key_size{ - opensslpp::cipher_context::get_key_size_in_bytes(active_data_cipher_)}; - - // generating random file key - util::hex_value_storage file_key{file_key_size}; - opensslpp::crypto_rng::generate(file_key); - - // creating an encryption context with the KEK cipher, the KEK, and - // the IV for file key encryption - opensslpp::cipher_context file_key_encryption_context{ - opensslpp::cipher_context_operation_type::encryption, kek_cipher, kek, - iv_for_file_key_encryption_v}; - - // identify the size of the file key encryption tag from the encryption - // (should be non-zero only for GCM modes) - const auto file_key_encryption_tag_size{ - file_key_encryption_context.get_tag_size_in_bytes()}; - // provisioning the optional storage for the file key encryption tag - util::optional_hex_value_storage tag_of_file_key_encryption{}; - util::byte_span tag_of_file_key_encryption_v{}; - if (file_key_encryption_tag_size != 0U) { - tag_of_file_key_encryption.emplace(file_key_encryption_tag_size); - tag_of_file_key_encryption_v = *tag_of_file_key_encryption; - } - - // performing the file key encryption and finalizing the tag (if any) - util::hex_value_storage file_key_encrypted_with_kek{file_key_size}; - file_key_encryption_context.update(file_key, file_key_encrypted_with_kek); - file_key_encryption_context.finalize(tag_of_file_key_encryption_v); - - // identifying the size of the IV that will be used for data encryption based - // on the active data cipher - const auto iv_length_for_data_encryption{ - opensslpp::cipher_context::get_iv_size_in_bytes(active_data_cipher_)}; - // generating random IV for file data encryption - util::hex_value_storage iv_for_data_encryption{iv_length_for_data_encryption}; - opensslpp::crypto_rng::generate(iv_for_data_encryption); - - // the tag of data encryption will be generated during the actual data - // encryption - binlog_encryption_record encryption_record{ - .kek_id = active_kek_id_, - .file_key_encrypted_with_kek = file_key_encrypted_with_kek, - .iv_for_file_key_encryption = iv_for_file_key_encryption, - .tag_of_file_key_encryption = tag_of_file_key_encryption, - .data_cipher = active_data_cipher_, - .iv_for_data_encryption = iv_for_data_encryption, - .tag_of_data_encryption = {}}; - - return encryption_record; -} - -void storage::write_data_to_stream( - util::const_byte_span data, - const optional_binlog_encryption_record &encryption_record, - std::uint64_t offset) { - if (!encryption_record.has_value()) { - // an early return when no encryption is needed - backend_->write_data_to_stream(data); - return; - } - - // as for security reasons our intent is to not store file keys in plaintext - // permanently, we need to decrypt the file key with the KEK before we can - // use it for data encryption. - - const auto &keyring_record{keyring_->get_key(encryption_record->kek_id)}; - - const auto &kek_cipher{keyring_record.get<"cipher">()}; - const auto &kek{keyring_record.get<"data_hex">().get_data()}; - - util::const_byte_span iv_for_file_key_encryption_v{}; - if (encryption_record->iv_for_file_key_encryption.has_value()) { - iv_for_file_key_encryption_v = - *encryption_record->iv_for_file_key_encryption; - }; - util::const_byte_span tag_of_file_key_encryption_v{}; - if (encryption_record->tag_of_file_key_encryption.has_value()) { - tag_of_file_key_encryption_v = - *encryption_record->tag_of_file_key_encryption; - } - - // creating a context for the file key decryption - opensslpp::cipher_context file_key_decryption_context{ - opensslpp::cipher_context_operation_type::decryption, kek_cipher, kek, - iv_for_file_key_encryption_v, tag_of_file_key_encryption_v}; - util::hex_value_storage file_key_decrypted{ - std::size(encryption_record->file_key_encrypted_with_kek)}; - file_key_decryption_context.update( - encryption_record->file_key_encrypted_with_kek, file_key_decrypted); - file_key_decryption_context.finalize(); - - // creating an context for data encryption with the data cipher, the file - // key (decrypted previously), and the IV for data encryption - - auto data_encryption_context{opensslpp::cipher_context::create_with_offset( - offset, opensslpp::cipher_context_operation_type::encryption, - encryption_record->data_cipher, file_key_decrypted, - encryption_record->iv_for_data_encryption)}; - - util::optional_hex_value_storage tag_of_data_encryption{}; - util::byte_span tag_of_data_encryption_v{}; - const auto data_encryption_tag_size{ - data_encryption_context.get_tag_size_in_bytes()}; - if (data_encryption_tag_size != 0U) { - tag_of_data_encryption.emplace(data_encryption_tag_size); - tag_of_data_encryption_v = *tag_of_data_encryption; + "operation requires storage to be constructed in streaming mode"); } - - util::hex_value_storage encrypted_data{std::size(data)}; - data_encryption_context.update(data, encrypted_data); - data_encryption_context.finalize(tag_of_data_encryption_v); - - backend_->write_data_to_stream(encrypted_data); - - // TODO: update file data encryption tag here, if one day we decide to - // support GCM mode for file data encryption } } // namespace binsrv diff --git a/src/binsrv/storage.hpp b/src/binsrv/storage.hpp index b274071..ede931a 100644 --- a/src/binsrv/storage.hpp +++ b/src/binsrv/storage.hpp @@ -31,6 +31,7 @@ #include "binsrv/encryption_format_type_fwd.hpp" #include "binsrv/main_config_fwd.hpp" #include "binsrv/replication_mode_type_fwd.hpp" +#include "binsrv/storage_core_fwd.hpp" #include "binsrv/events/composite_binlog_name.hpp" @@ -50,46 +51,6 @@ namespace binsrv { class [[nodiscard]] storage { public: - struct binlog_encryption_record { - std::string kek_id; - util::hex_value_storage file_key_encrypted_with_kek; - util::optional_hex_value_storage iv_for_file_key_encryption; - util::optional_hex_value_storage tag_of_file_key_encryption; - std::string data_cipher; - util::hex_value_storage iv_for_data_encryption; - util::optional_hex_value_storage tag_of_data_encryption; - - [[nodiscard]] static models::binlog_file_encryption_record - to_model(const binlog_encryption_record &record); - [[nodiscard]] static binlog_encryption_record - from_model(const models::binlog_file_encryption_record &model); - }; - using optional_binlog_encryption_record = - std::optional; - struct binlog_record { - // binlog file name - events::composite_binlog_name name; - // binlog file size in bytes - std::uint64_t size{0ULL}; - // accumulated GTIDs present in the binlog files before this one - gtids::optional_gtid_set previous_gtids{}; - // GTIDs present in this binlog file - gtids::optional_gtid_set added_gtids{}; - // minimum and maximum event timestamps observed in this binlog file - util::ctime_timestamp_range timestamps{}; - // sequence_number of the last transaction seen in this file - - // used for GTID rewrite-mode resume state persistence - events::seq_no_t last_sequence_number{0ULL}; - // optional encryption parameters - optional_binlog_encryption_record encryption{}; - }; - using binlog_record_container = std::vector; - - static constexpr std::string_view default_binlog_index_name{"binlog.index"}; - static constexpr std::string_view default_binlog_index_entry_path{"."}; - static constexpr std::string_view metadata_name{"metadata.json"}; - static constexpr std::string_view binlog_metadata_extension{".json"}; - static constexpr std::size_t default_event_buffer_size_in_bytes{16384U}; storage(basic_logger_ptr logger, const main_config &config, @@ -100,58 +61,26 @@ class [[nodiscard]] storage { storage(storage &&) = delete; storage &operator=(storage &&) = delete; - // destructor is explicitly declared here and defined as default in .cpp - // file to complete the rule of 5 ~storage(); - [[nodiscard]] const gtids::gtid_set &get_purged_gtids() const noexcept { - return purged_gtids_; - } + [[nodiscard]] const gtids::gtid_set &get_purged_gtids() const noexcept; 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 { - return replication_mode_; - } + [[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 { - return binlog_records_; - } - [[nodiscard]] bool is_empty() const noexcept { - return binlog_records_.empty(); - } - [[nodiscard]] const events::composite_binlog_name & - get_current_binlog_name() const noexcept { - return is_empty() ? binlog_name_sentinel_ - : get_current_binlog_record().name; - } + get_binlog_records() const noexcept; + [[nodiscard]] bool is_empty() const noexcept; + [[nodiscard]] events::composite_binlog_name get_current_binlog_name() const; + [[nodiscard]] std::uint64_t get_current_position() const noexcept { return get_flushed_position() + std::size(event_buffer_); } - [[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; - } + [[nodiscard]] gtids::gtid_set get_gtids() const; [[nodiscard]] events::seq_no_t get_last_transaction_sequence_number() const noexcept { @@ -192,30 +121,15 @@ class [[nodiscard]] storage { [[nodiscard]] std::string get_binlog_uri(const events::composite_binlog_name &binlog_name) const; - [[nodiscard]] bool is_keyring_initialized() const noexcept { - return static_cast(keyring_); - } + [[nodiscard]] bool is_keyring_initialized() const noexcept; [[nodiscard]] std::string get_keyring_description() const; [[nodiscard]] std::string get_active_kek_description() const; [[nodiscard]] std::string get_encryption_format_description() const; - [[nodiscard]] bool has_active_kek() const noexcept { - return !active_kek_id_.empty(); - } + [[nodiscard]] bool has_active_kek() const noexcept; private: - basic_logger_ptr logger_; - storage_construction_mode_type construction_mode_; - basic_keyring_ptr keyring_; - optional_encryption_format_type encryption_format_; - std::string active_kek_id_; - std::string active_data_cipher_{}; - basic_storage_backend_ptr backend_; - - replication_mode_type replication_mode_; - events::composite_binlog_name binlog_name_sentinel_{}; - gtids::gtid_set purged_gtids_{}; - binlog_record_container binlog_records_{}; + storage_core_ptr core_; std::uint64_t checkpoint_size_bytes_{0ULL}; std::uint64_t last_checkpoint_position_{0ULL}; @@ -232,22 +146,6 @@ class [[nodiscard]] storage { events::seq_no_t ready_to_flush_last_sequence_number_{0ULL}; events::seq_no_t incomplete_transaction_last_sequence_number_{0ULL}; - void remove_temporary_objects(storage_object_name_container &object_names); - - void initialize_storage_encryption( - const optional_encryption_config &encryption_config); - - void ensure_streaming_mode() const; - void ensure_purging_mode() const; - - [[nodiscard]] const binlog_record & - get_current_binlog_record() const noexcept { - return binlog_records_.back(); - } - [[nodiscard]] binlog_record &get_current_binlog_record() noexcept { - return binlog_records_.back(); - } - [[nodiscard]] bool size_checkpointing_enabled() const noexcept { return checkpoint_size_bytes_ != 0ULL; } @@ -261,9 +159,7 @@ 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 { - return is_empty() ? 0ULL : get_current_binlog_record().size; - } + [[nodiscard]] std::uint64_t get_flushed_position() const noexcept; [[nodiscard]] std::uint64_t get_ready_to_flush_position() const noexcept { return get_flushed_position() + last_transaction_boundary_position_in_event_buffer_; @@ -275,35 +171,7 @@ class [[nodiscard]] storage { void flush_event_buffer_internal(); - void load_binlog_index(); - void validate_binlog_index( - const storage_object_name_container &object_names) const; - void save_binlog_index() const; - - void load_metadata(); - void validate_metadata( - replication_mode_type replication_mode, - const optional_encryption_format_type &encryption_format) const; - void save_metadata() const; - - [[nodiscard]] static std::string generate_binlog_metadata_name( - const events::composite_binlog_name &binlog_name); - [[nodiscard]] binlog_record - load_binlog_metadata(const events::composite_binlog_name &binlog_name) const; - void validate_binlog_metadata(const binlog_record &record) const; - void save_binlog_metadata(const binlog_record &record) const; - - void load_and_validate_binlog_metadata_set( - const storage_object_name_container &object_names, - const storage_object_name_container &object_metadata_names); - - [[nodiscard]] optional_binlog_encryption_record - generate_binlog_encryption_record() const; - - void write_data_to_stream( - util::const_byte_span data, - const optional_binlog_encryption_record &encryption_record, - std::uint64_t offset); + void ensure_streaming_mode() const; }; } // namespace binsrv diff --git a/src/binsrv/storage_core.cpp b/src/binsrv/storage_core.cpp new file mode 100644 index 0000000..33edd8e --- /dev/null +++ b/src/binsrv/storage_core.cpp @@ -0,0 +1,987 @@ +// Copyright (c) 2023-2024 Percona and/or its affiliates. +// +// This program is free software; you can redistribute it and/or modify +// it under the terms of the GNU General Public License, version 2.0, +// as published by the Free Software Foundation. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU General Public License, version 2.0, for more details. +// +// You should have received a copy of the GNU General Public License +// along with this program; if not, write to the Free Software +// Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA + +#include "binsrv/storage_core.hpp" + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "binsrv/basic_keyring.hpp" +#include "binsrv/basic_logger.hpp" +#include "binsrv/basic_storage_backend.hpp" +#include "binsrv/binlog_file_metadata.hpp" +#include "binsrv/encryption_format_type.hpp" +#include "binsrv/keyring_factory.hpp" +#include "binsrv/keyring_record.hpp" +#include "binsrv/log_severity.hpp" +#include "binsrv/main_config.hpp" +#include "binsrv/replication_config.hpp" +#include "binsrv/replication_mode_type.hpp" +#include "binsrv/storage_backend_factory.hpp" +#include "binsrv/storage_config.hpp" +#include "binsrv/storage_metadata.hpp" + +#include "binsrv/events/common_types.hpp" +#include "binsrv/events/composite_binlog_name.hpp" +#include "binsrv/events/protocol_traits_fwd.hpp" + +#include "binsrv/gtids/gtid_set.hpp" + +#include "binsrv/models/binlog_file_encryption_record.hpp" + +#include "opensslpp/cipher_context.hpp" +#include "opensslpp/crypto_rng.hpp" + +#include "util/byte_span.hpp" +#include "util/ctime_timestamp_range.hpp" +#include "util/exception_location_helpers.hpp" + +namespace binsrv { + +[[nodiscard]] models::binlog_file_encryption_record +binlog_encryption_record::to_model(const binlog_encryption_record &record) { + models::binlog_file_encryption_record model{}; + + auto &file_key_envelope{model.get<"file_key_envelope">()}; + file_key_envelope.get<"kek_id">() = record.kek_id; + file_key_envelope.get<"data_hex">() = record.file_key_encrypted_with_kek; + if (record.iv_for_file_key_encryption.has_value()) { + file_key_envelope.get<"iv_hex">() = *record.iv_for_file_key_encryption; + } + if (record.tag_of_file_key_encryption.has_value()) { + file_key_envelope.get<"tag_hex">() = *record.tag_of_file_key_encryption; + } + + auto &file_data_envelope{model.get<"file_data_envelope">()}; + file_data_envelope.get<"cipher">() = record.data_cipher; + file_data_envelope.get<"iv_hex">() = record.iv_for_data_encryption; + if (record.tag_of_data_encryption.has_value()) { + file_data_envelope.get<"tag_hex">() = *record.tag_of_data_encryption; + } + return model; +} + +[[nodiscard]] binlog_encryption_record binlog_encryption_record::from_model( + const models::binlog_file_encryption_record &model) { + binlog_encryption_record record{}; + + const auto &file_key_envelope{model.get<"file_key_envelope">()}; + record.kek_id = file_key_envelope.get<"kek_id">(); + const auto file_key_raw{file_key_envelope.get<"data_hex">().get_data()}; + record.file_key_encrypted_with_kek.assign(std::cbegin(file_key_raw), + std::cend(file_key_raw)); + if (file_key_envelope.get<"iv_hex">().has_value()) { + const auto file_key_iv_raw{file_key_envelope.get<"iv_hex">()->get_data()}; + record.iv_for_file_key_encryption.emplace(std::cbegin(file_key_iv_raw), + std::cend(file_key_iv_raw)); + } + if (file_key_envelope.get<"tag_hex">().has_value()) { + const auto file_key_tag_raw{file_key_envelope.get<"tag_hex">()->get_data()}; + record.tag_of_file_key_encryption.emplace(std::cbegin(file_key_tag_raw), + std::cend(file_key_tag_raw)); + } + + const auto &file_data_envelope{model.get<"file_data_envelope">()}; + record.data_cipher = file_data_envelope.get<"cipher">(); + const auto file_data_iv_raw{file_data_envelope.get<"iv_hex">().get_data()}; + record.iv_for_data_encryption.assign(std::cbegin(file_data_iv_raw), + std::cend(file_data_iv_raw)); + if (file_data_envelope.get<"tag_hex">().has_value()) { + const auto file_data_tag_raw{ + file_data_envelope.get<"tag_hex">()->get_data()}; + record.tag_of_data_encryption.emplace(std::cbegin(file_data_tag_raw), + std::cend(file_data_tag_raw)); + } + return record; +} + +storage_core::storage_core(basic_logger_ptr logger, const main_config &config, + storage_construction_mode_type construction_mode) + : logger_{std::move(logger)}, construction_mode_{construction_mode}, + backend_{} { + assert(logger_); + + const auto &replication_config{config.root().get<"replication">()}; + // we need a copy of replication mode as replication_mode_ will be + // overwritten by load_metadata() later + const auto replication_mode{replication_config.get<"mode">()}; + replication_mode_ = replication_mode; + + const auto &storage_config{config.root().get<"storage">()}; + + const auto &keyring_config{config.root().get<"keyring">()}; + if (keyring_config.has_value()) { + keyring_ = keyring_factory::create(keyring_config->get<"uri">()); + } + const auto &encryption_config{storage_config.get<"encryption">()}; + initialize_storage_encryption(encryption_config); + + backend_ = storage_backend_factory::create(storage_config); + + auto storage_objects{backend_->list_objects()}; + remove_temporary_objects(storage_objects); + + if (storage_objects.empty()) { + // initialized on a new / empty storage - just save metadata and return + if (construction_mode_ == storage_construction_mode_type::streaming) { + save_metadata(); + } + return; + } + + const auto metadata_it{std::as_const(storage_objects).find(metadata_name)}; + if (metadata_it == std::cend(storage_objects)) { + util::exception_location().raise( + "storage is not empty but does not contain metadata"); + } + storage_objects.erase(metadata_it); + + // as load_metadata() will be updating 'encryption_format_', saving it here + // to use for validation later + const auto encryption_format{encryption_format_}; + load_metadata(); + validate_metadata(replication_mode, encryption_format); + // in case when storage metadata file is present and did not have encryption + // format specified, but it is set in the configuration file, we need to + // update storage metadata + if (!encryption_format_.has_value() && encryption_format.has_value()) { + encryption_format_ = encryption_format; + save_metadata(); + } + + // if after metadata erasure 'storage_objects' is empty, then this means + // that it has only metadata in it that passes validation and we can + // consider it as an initialized empty storage, so just return + if (storage_objects.empty()) { + return; + } + + const auto binlog_index_it{storage_objects.find(default_binlog_index_name)}; + if (binlog_index_it == std::cend(storage_objects)) { + // as binlog index file is created after the very first binlog data file + // and its metadata file are created, the concurrent query-only operation + // may intrude exactly between these two steps, so for query-only operations + // we should not consider the absence of binlog index file as an error + if (construction_mode == storage_construction_mode_type::querying_only) { + return; + } + util::exception_location().raise( + "storage is not empty but does not contain binlog index"); + } + storage_objects.erase(binlog_index_it); + + // extracting all binlog file metadata files into a separate container + storage_object_name_container storage_metadata_objects; + for (auto storage_object_it{std::cbegin(storage_objects)}; + storage_object_it != std::cend(storage_objects);) { + const std::filesystem::path object_name{storage_object_it->first}; + if (object_name.has_extension() && + object_name.extension() == binlog_metadata_extension) { + auto object_node = storage_objects.extract(storage_object_it++); + storage_metadata_objects.insert(std::move(object_node)); + } else { + ++storage_object_it; + } + } + load_binlog_index(); + validate_binlog_index(storage_objects); + + load_and_validate_binlog_metadata_set(storage_objects, + storage_metadata_objects); + assert(!binlog_records_.front().added_gtids.has_value() || + purged_gtids_ == binlog_records_.front().added_gtids); +} + +storage_core::~storage_core() = default; + +void storage_core::set_purged_gtids(const gtids::gtid_set &purged_gtids) { + if (!is_in_gtid_replication_mode()) { + util::exception_location().raise( + "cannot set purged GTIDs in position-based replication mode"); + } + if (!is_empty()) { + util::exception_location().raise( + "cannot set purged GTIDs in a non-empty storage"); + } + purged_gtids_ = purged_gtids; +} + +[[nodiscard]] std::string storage_core::get_backend_description() const { + return backend_->get_description(); +} + +[[nodiscard]] bool storage_core::is_in_gtid_replication_mode() const noexcept { + return replication_mode_ == replication_mode_type::gtid; +} + +[[nodiscard]] bool storage_core::is_binlog_open() const noexcept { + return backend_->is_stream_open(); +} + +[[nodiscard]] open_binlog_status +storage_core::open_binlog(const events::composite_binlog_name &binlog_name) { + ensure_streaming_mode(); + + auto result{open_binlog_status::opened_with_data_present}; + + // here we either create a new binlog file if its name is not presentin the + // "binlog_records_", or we open an existing one and append to it, in which + // case we need to make sure that the current position is properly set + const bool binlog_exists{ + std::ranges::find(std::as_const(binlog_records_), binlog_name, + &binlog_record::name) != std::cend(binlog_records_)}; + + // 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()) { + util::exception_location().raise( + "cannot open an existing binlog that is not the latest one for " + "append"); + } + if (get_flushed_position() == 0ULL) { + util::exception_location().raise( + "invalid position set when opening an existing binlog"); + } + } + + const auto mode{binlog_exists ? storage_backend_open_stream_mode::append + : storage_backend_open_stream_mode::create}; + const auto open_stream_offset{backend_->open_stream(binlog_name.str(), mode)}; + + if (binlog_exists) { + result = open_existing_binlog_file_internal(open_stream_offset); + } else { + result = open_new_binlog_file_internal(binlog_name); + } + + return result; +} + +void storage_core::write_event_block( + util::const_byte_span event_block_data, const gtids::gtid_set &block_gtids, + const util::ctime_timestamp_range &block_timestamps, + 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); + if (is_in_gtid_replication_mode()) { + auto &optional_added_gtids{get_current_binlog_record().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; + + save_binlog_metadata(get_current_binlog_record()); +} + +void storage_core::close_binlog() { + ensure_streaming_mode(); + + backend_->close_stream(); +} + +[[nodiscard]] std::pair +storage_core::purge_binlogs(const events::composite_binlog_name &target) { + ensure_purging_mode(); + + if (is_empty()) { + util::exception_location().raise( + "cannot purge: binlog storage is empty"); + } + const auto &front_base_name{binlog_records_.front().name.get_base_name()}; + if (target.get_base_name() != front_base_name) { + util::exception_location().raise( + "cannot purge: target binlog name has a different base name than " + "the binlog records in the storage"); + } + const auto target_it{std::ranges::find(std::as_const(binlog_records_), target, + &binlog_record::name)}; + if (target_it == std::cend(binlog_records_)) { + util::exception_location().raise( + "cannot purge: target binlog name is not present in the storage"); + } + // refuse to purge the current tail: emptying the storage would lose + // the resume position (current binlog name / position in position + // mode, executed GTID set in GTID mode) and force the next 'fetch' / + // 'pull' to re-stream from the very beginning of the source's + // retained binlog history. + if (target_it == std::prev(std::cend(binlog_records_))) { + util::exception_location().raise( + "cannot purge: target is the current tail binlog file; at least " + "one binlog file must remain in the storage to preserve the " + "resume position"); + } + + // step 1: extract the prefix [begin, target_it + 1) - this + // becomes the set of records we are going to drop on disk; the + // returned vector preserves the original order so the caller can + // use it directly to produce a response + const auto victim_count{static_cast( + std::distance(std::cbegin(binlog_records_), target_it) + 1)}; + binlog_record_container removed_records; + removed_records.reserve(victim_count); + std::move(std::begin(binlog_records_), + std::begin(binlog_records_) + + static_cast(victim_count), + std::back_inserter(removed_records)); + binlog_records_.erase(std::begin(binlog_records_), + std::begin(binlog_records_) + + static_cast(victim_count)); + + // step 2: rewrite the binlog index from the surviving records left + // in 'binlog_records_' after step 1 (always non-empty thanks to the + // tail-refusal guard above). 'save_binlog_index' goes through the + // backend's atomic-overwrite 'put_object', so from this point on + // the purge is considered committed - any subsequent failure + // leaves the storage in an inconsistent state (leftover payload / + // metadata files no longer referenced by the index) that the + // constructor's existing validators will refuse to open on next + // startup. + save_binlog_index(); + + // step 3: best-effort removal of the victim payload + metadata + // objects; any failure here is intentionally swallowed - the index + // has already been committed and reporting a "file could not be + // removed" error to the caller would falsely suggest that the + // purge itself failed; the resulting leftovers will trip the + // constructor's validators on next startup. + // We materialise the (metadata + payload) names for every victim + // into a single batch and hand it to 'basic_storage_backend:: + // remove_objects', which runs the backend's durability barrier + // exactly once at the end of the batch - so the whole batch + // amortises to a single fsync(2) on the local filesystem backend + // (and a no-op on S3) instead of O(N) syncs. + std::vector victim_object_names; + victim_object_names.reserve(std::size(removed_records) * 2U); + for (const auto &victim : removed_records) { + victim_object_names.emplace_back( + generate_binlog_metadata_name(victim.name)); + victim_object_names.emplace_back(victim.name.str()); + } + std::string cleanup_warning_message; + try { + backend_->remove_objects(victim_object_names); + } catch (const std::exception &e) { + // 'remove_objects' re-raises the first per-name failure (if any) + // after running the durability barrier; we do not propagate it + // to the caller because the index has already been committed + // and any leftover payload/metadata files will be picked up by + // the constructor's validators on the next startup. We just + // capture the underlying message so the caller can surface it + // under a 'warning' status in the JSON response. + cleanup_warning_message = e.what(); + } + + return {std::move(removed_records), std::move(cleanup_warning_message)}; +} + +[[nodiscard]] std::string storage_core::get_binlog_uri( + const events::composite_binlog_name &binlog_name) const { + return backend_->get_object_uri(binlog_name.str()); +} + +[[nodiscard]] std::string storage_core::get_keyring_description() const { + return is_keyring_initialized() ? keyring_->get_description() + : "keyring is not initialized"; +} + +[[nodiscard]] std::string storage_core::get_active_kek_description() const { + 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 { + return encryption_format_.has_value() + ? std::string{to_string_view(*encryption_format_)} + : std::string{"encryption format is not set"}; +} + +void storage_core::remove_temporary_objects( + storage_object_name_container &object_names) { + using remove_object_container = std::vector; + remove_object_container remove_objects; + for (auto it{std::begin(object_names)}; it != std::end(object_names);) { + const std::filesystem::path object_name{it->first}; + if (object_name.has_extension() && + object_name.extension() == tmp_storage_object_suffix) { + auto object_node = object_names.extract(it++); + remove_objects.emplace_back(std::move(object_node.key())); + logger_->log_format(log_severity::warning, + "found temporary storage object '{}' left after " + "improper shutdown", + object_name.string()); + } else { + ++it; + } + } + if (remove_objects.empty()) { + return; + } + + // for querying-only mode we do not perform any modifying operations - just + // clear the 'object_names' container and return + if (construction_mode_ == storage_construction_mode_type::querying_only) { + return; + } + + backend_->remove_objects(remove_objects); + logger_->log_format(log_severity::warning, + "removed {} temporary storage object(s) left after " + "improper shutdown", + remove_objects.size()); +} + +void storage_core::initialize_storage_encryption( + const optional_encryption_config &encryption_config) { + if (encryption_config.has_value()) { + encryption_format_ = encryption_config->get<"format">(); + if (!is_keyring_initialized()) { + util::exception_location().raise( + "encryption is enabled but keyring is not initialized"); + } + const auto &kek_id{encryption_config->get<"kek_id">()}; + if (!keyring_->contains(kek_id)) { + util::exception_location().raise( + "keyring does not contain the specified KEK ID"); + } + active_kek_id_ = kek_id; + active_data_cipher_ = encryption_config->get<"cipher">(); + + // make sure that random file keys (of length that corresponds to the + // active data cipher) can be encrypted with the active KEK - + // for instance, if active data cipher is AES-192-CRT (key length 24 + // bytes), then the active KEK cannot be of ECB or CBC mode as these + // ciphers can only encrypt data of length that is a multiple of the + // block size (16 bytes) + const auto &keyring_record{keyring_->get_key(active_kek_id_)}; + if (opensslpp::cipher_context::get_key_size_in_bytes(active_data_cipher_) % + opensslpp::cipher_context::get_block_size_in_bytes( + keyring_record.get<"cipher">()) != + 0U) { + util::exception_location().raise( + "active data cipher key length is not compatible with the active " + "KEK cipher block size"); + } + } +} + +void storage_core::ensure_streaming_mode() const { + if (construction_mode_ != storage_construction_mode_type::streaming) { + util::exception_location().raise( + "operation requires storage to be constructed in streaming mode"); + } +} + +void storage_core::ensure_purging_mode() const { + if (construction_mode_ != storage_construction_mode_type::purging) { + util::exception_location().raise( + "operation requires storage to be constructed in purging mode"); + } +} + +[[nodiscard]] open_binlog_status storage_core::open_new_binlog_file_internal( + const events::composite_binlog_name &binlog_name) { + + auto encryption_record{generate_binlog_encryption_record()}; + // writing the magic binlog footprint only if this is a newly + // created file + write_data_to_stream(events::magic_binlog_payload, encryption_record, 0ULL); + + 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(); + added_binlog_gtids = gtids::gtid_set{}; + } + + binlog_records_.emplace_back( + binlog_name, events::magic_binlog_offset, + 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_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); + if (open_stream_offset >= events::magic_binlog_offset) { + return open_stream_offset == events::magic_binlog_offset + ? open_binlog_status::opened_at_magic_payload_offset + : open_binlog_status::opened_with_data_present; + } + 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; + return open_binlog_status::opened_empty; +} + +void storage_core::load_binlog_index() { + const auto index_content{backend_->get_object(default_binlog_index_name)}; + // opening in text mode + std::istringstream index_iss{index_content}; + std::string current_line; + while (std::getline(index_iss, current_line)) { + if (current_line.empty()) { + continue; + } + const std::filesystem::path current_binlog_path{current_line}; + if (current_binlog_path.parent_path() != default_binlog_index_entry_path) { + util::exception_location().raise( + "binlog index contains an entry that has an invalid path"); + } + auto current_binlog_name{current_binlog_path.filename().string()}; + + if (current_binlog_name == default_binlog_index_name) { + util::exception_location().raise( + "binlog index contains a reference to the binlog index name"); + } + const auto current_binlog_name_parsed{ + events::composite_binlog_name::parse(current_binlog_name)}; + if (std::ranges::find(std::as_const(binlog_records_), + current_binlog_name_parsed, + &binlog_record::name) != std::cend(binlog_records_)) { + util::exception_location().raise( + "binlog index contains a duplicate entry"); + } + gtids::optional_gtid_set previous_binlog_gtids{}; + gtids::optional_gtid_set added_binlog_gtids{}; + if (is_in_gtid_replication_mode()) { + previous_binlog_gtids = gtids::gtid_set{}; + added_binlog_gtids = gtids::gtid_set{}; + } + binlog_records_.emplace_back( + current_binlog_name_parsed, 0ULL, std::move(previous_binlog_gtids), + std::move(added_binlog_gtids), util::ctime_timestamp_range{}); + } +} + +void storage_core::validate_binlog_index( + const storage_object_name_container &object_names) const { + // in the querying_only mode we allow discrepancies between the binlog index + // and the actual objects in the storage + if (construction_mode_ == storage_construction_mode_type::querying_only) { + return; + } + + for (auto const &record : binlog_records_) { + if (!object_names.contains(record.name.str())) { + util::exception_location().raise( + "binlog index contains a reference to a non-existing object"); + } + } + + if (std::size(object_names) != std::size(binlog_records_)) { + util::exception_location().raise( + "storage contains an object that is not " + "referenced in the binlog index"); + } + + // TODO: add integrity checks (parsing + checksumming) for the binlog + // files in the index +} + +void storage_core::save_binlog_index() const { + std::ostringstream oss; + for (const auto &record : binlog_records_) { + std::filesystem::path binlog_path{default_binlog_index_entry_path}; + binlog_path /= record.name.str(); + oss << binlog_path.generic_string() << '\n'; + } + const auto content{oss.str()}; + backend_->put_object(default_binlog_index_name, + util::as_const_byte_span(content)); +} + +void storage_core::load_metadata() { + const auto metadata_content{backend_->get_object(metadata_name)}; + const storage_metadata metadata{metadata_content}; + replication_mode_ = metadata.root().get<"mode">(); + encryption_format_ = metadata.root().get<"encryption">(); +} + +void storage_core::validate_metadata( + replication_mode_type replication_mode, + const optional_encryption_format_type &encryption_format) const { + if (replication_mode != replication_mode_) { + util::exception_location().raise( + "replication mode provided to initialize storage differs from the one " + "stored in metadata"); + } + + if (encryption_format_.has_value() && encryption_format.has_value() && + *encryption_format_ != *encryption_format) { + // if both encryption formats (the existing one loaded from the storage + // metadata file and a new one specified in the configuration file) are + // present and have different values, then this is an error + util::exception_location().raise( + "storage encryption format provided to initialize storage differs from " + "the one stored in metadata"); + } +} + +void storage_core::save_metadata() const { + storage_metadata metadata{}; + metadata.root().get<"mode">() = replication_mode_; + metadata.root().get<"encryption">() = encryption_format_; + const auto content{metadata.str()}; + backend_->put_object(metadata_name, util::as_const_byte_span(content)); +} + +[[nodiscard]] std::string storage_core::generate_binlog_metadata_name( + const events::composite_binlog_name &binlog_name) { + std::string binlog_metadata_name{binlog_name.str()}; + binlog_metadata_name += storage_core::binlog_metadata_extension; + return binlog_metadata_name; +} + +[[nodiscard]] binlog_record storage_core::load_binlog_metadata( + const events::composite_binlog_name &binlog_name) const { + const auto content{ + backend_->get_object(generate_binlog_metadata_name(binlog_name))}; + binlog_file_metadata metadata{content}; + + const auto &optional_encryption_metadata{metadata.root().get<"encryption">()}; + return binlog_record{ + .name = binlog_name, + .size = metadata.root().get<"size">(), + .previous_gtids = metadata.root().get<"previous_gtids">(), + .added_gtids = metadata.root().get<"added_gtids">(), + .timestamps = {metadata.root().get<"min_timestamp">(), + metadata.root().get<"max_timestamp">()}, + .last_sequence_number = metadata.root().get<"last_sequence_number">(), + .encryption = optional_encryption_metadata.has_value() + ? binlog_encryption_record::from_model( + *optional_encryption_metadata) + : optional_binlog_encryption_record{}}; +} + +void storage_core::validate_binlog_metadata(const binlog_record &record) const { + if (is_in_gtid_replication_mode()) { + if (!record.previous_gtids.has_value()) { + util::exception_location().raise( + "missing previous GTID set in the binlog metadata while in GTID " + "replication " + "mode"); + } + if (!record.added_gtids.has_value()) { + util::exception_location().raise( + "missing added GTID set in the binlog metadata while in GTID " + "replication " + "mode"); + } + } else { + if (record.previous_gtids.has_value()) { + util::exception_location().raise( + "found previous GTID set in the binlog metadata while in position " + "replication mode"); + } + if (record.added_gtids.has_value()) { + util::exception_location().raise( + "found added GTID set in the binlog metadata while in position " + "replication mode"); + } + } + // make sure that if the encryption record is present in the binlog + // metadata, keyring must be initialized and contain the KEK with the + // ID specified in the encryption record + if (record.encryption.has_value()) { + if (!is_keyring_initialized()) { + util::exception_location().raise( + "found encryption record in the binlog metadata but keyring is not " + "initialized"); + } + if (!keyring_->contains(record.encryption->kek_id)) { + util::exception_location().raise( + "found encryption record in the binlog metadata but keyring does not " + "contain the specified KEK ID"); + } + } +} + +void storage_core::save_binlog_metadata(const binlog_record &record) const { + binlog_file_metadata metadata{}; + metadata.root().get<"size">() = record.size; + metadata.root().get<"previous_gtids">() = record.previous_gtids; + metadata.root().get<"added_gtids">() = record.added_gtids; + metadata.root().get<"min_timestamp">() = + util::ctime_timestamp{record.timestamps.get_min_timestamp()}; + metadata.root().get<"max_timestamp">() = + util::ctime_timestamp{record.timestamps.get_max_timestamp()}; + metadata.root().get<"last_sequence_number">() = record.last_sequence_number; + const auto &record_encryption{record.encryption}; + if (record_encryption.has_value()) { + metadata.root().get<"encryption">() = + binlog_encryption_record::to_model(*record_encryption); + } + const auto content{metadata.str()}; + backend_->put_object(generate_binlog_metadata_name(record.name), + util::as_const_byte_span(content)); +} + +void storage_core::load_and_validate_binlog_metadata_set( + // NOLINTNEXTLINE(bugprone-easily-swappable-parameters) + const storage_object_name_container &object_names, + const storage_object_name_container &object_metadata_names) { + auto record_it{std::begin(binlog_records_)}; + while (record_it != std::end(binlog_records_)) { + binlog_record loaded_binlog_metadata{}; + try { + const auto binlog_metadata_name{ + generate_binlog_metadata_name(record_it->name)}; + if (!object_metadata_names.contains(binlog_metadata_name)) { + util::exception_location().raise( + "missing metadata for a binlog listed in the binlog index"); + } + loaded_binlog_metadata = load_binlog_metadata(record_it->name); + } catch (const std::exception &) { + if (construction_mode_ == storage_construction_mode_type::querying_only) { + // in the querying_only mode we just skip invalid metadata and the + // corresponding binlog file - this allows to query an otherwise + // unusable storage and retrieve information about valid binlog files + // from it, which can be useful for debugging / forensics purposes + record_it = binlog_records_.erase(record_it); + continue; + } + throw; + } + validate_binlog_metadata(loaded_binlog_metadata); + // validating binlog size from the metadata only makes sense if we are not + // in the querying_only mode + if (construction_mode_ != storage_construction_mode_type::querying_only) { + const auto binlog_file_name{record_it->name.str()}; + const auto actual_binlog_file_size{object_names.at(binlog_file_name)}; + // validating that the size stored in the metadata matches the actual size + if (loaded_binlog_metadata.size != actual_binlog_file_size) { + // in case when Binlog Server process was not properly shut down + // there is a chance that there will be mismatch between the actual + // binlog data file size and the 'size' field in the binlog metadata + + // if this mismatch was found in the most recent binlog file, we can + // perform automatic recovery (truncating binlog data file content to + // the size from the metadata) + if (std::next(record_it) != std::end(binlog_records_)) { + util::exception_location().raise( + "size from the binlog metadata does not match the actual binlog " + "size"); + } + // performing recovery + if (loaded_binlog_metadata.size > actual_binlog_file_size) { + util::exception_location().raise( + "cannot perform recovery - size from the binlog metadata is " + "bigger than the actual binlog file size"); + } + backend_->resize_object(binlog_file_name, loaded_binlog_metadata.size); + logger_->log_format(log_severity::warning, + "recovered binlog file '{}' by truncating it to " + "the size from the metadata", + binlog_file_name); + } + } + *record_it = std::move(loaded_binlog_metadata); + ++record_it; + } + // after this loop position_ and gtids_ should store the values from the last + // binlog file metadata + + if (construction_mode_ != storage_construction_mode_type::querying_only) { + if (std::size(object_metadata_names) != std::size(binlog_records_)) { + util::exception_location().raise( + "found metadata for a non-existing binlog"); + } + } + + // if we are in GTID replication mode, then we can consider GTIDs from the + // first binlog metadata as purged GTIDs for the whole storage + const auto &optional_added_gtids{binlog_records_.front().added_gtids}; + if (optional_added_gtids.has_value()) { + purged_gtids_ = *optional_added_gtids; + } +} + +[[nodiscard]] optional_binlog_encryption_record +storage_core::generate_binlog_encryption_record() const { + if (!has_active_kek()) { + return std::nullopt; + } + + // we identify the KEK record in the keyring by the active KEK ID, + // specified in the main configuration file + // ('' parameter) + const auto &keyring_record{keyring_->get_key(active_kek_id_)}; + + // identifying the the cipher name and the key data from the + // keyring record - this data will be used to encrypt random file + // keys generated for new binlog data files + const auto &kek_cipher{keyring_record.get<"cipher">()}; + const auto &kek{keyring_record.get<"data_hex">().get_data()}; + + // identify the size of the IV that will be used for file key + // encryption based on the KEK cipher; if the KEK cipher is in ECB mode, then + // the IV is not used and its size will be 0 + const auto iv_size_for_file_key_encryption{ + opensslpp::cipher_context::get_iv_size_in_bytes(kek_cipher)}; + util::optional_hex_value_storage iv_for_file_key_encryption{}; + util::const_byte_span iv_for_file_key_encryption_v{}; + if (iv_size_for_file_key_encryption != 0U) { + // generating random IV for file key encryption + iv_for_file_key_encryption.emplace(iv_size_for_file_key_encryption); + opensslpp::crypto_rng::generate(*iv_for_file_key_encryption); + iv_for_file_key_encryption_v = *iv_for_file_key_encryption; + } + + // identify the size of the file key based on the active data cipher + const auto file_key_size{ + opensslpp::cipher_context::get_key_size_in_bytes(active_data_cipher_)}; + + // generating random file key + util::hex_value_storage file_key{file_key_size}; + opensslpp::crypto_rng::generate(file_key); + + // creating an encryption context with the KEK cipher, the KEK, and + // the IV for file key encryption + opensslpp::cipher_context file_key_encryption_context{ + opensslpp::cipher_context_operation_type::encryption, kek_cipher, kek, + iv_for_file_key_encryption_v}; + + // identify the size of the file key encryption tag from the encryption + // (should be non-zero only for GCM modes) + const auto file_key_encryption_tag_size{ + file_key_encryption_context.get_tag_size_in_bytes()}; + // provisioning the optional storage for the file key encryption tag + util::optional_hex_value_storage tag_of_file_key_encryption{}; + util::byte_span tag_of_file_key_encryption_v{}; + if (file_key_encryption_tag_size != 0U) { + tag_of_file_key_encryption.emplace(file_key_encryption_tag_size); + tag_of_file_key_encryption_v = *tag_of_file_key_encryption; + } + + // performing the file key encryption and finalizing the tag (if any) + util::hex_value_storage file_key_encrypted_with_kek{file_key_size}; + file_key_encryption_context.update(file_key, file_key_encrypted_with_kek); + file_key_encryption_context.finalize(tag_of_file_key_encryption_v); + + // identifying the size of the IV that will be used for data encryption based + // on the active data cipher + const auto iv_length_for_data_encryption{ + opensslpp::cipher_context::get_iv_size_in_bytes(active_data_cipher_)}; + // generating random IV for file data encryption + util::hex_value_storage iv_for_data_encryption{iv_length_for_data_encryption}; + opensslpp::crypto_rng::generate(iv_for_data_encryption); + + // the tag of data encryption will be generated during the actual data + // encryption + binlog_encryption_record encryption_record{ + .kek_id = active_kek_id_, + .file_key_encrypted_with_kek = file_key_encrypted_with_kek, + .iv_for_file_key_encryption = iv_for_file_key_encryption, + .tag_of_file_key_encryption = tag_of_file_key_encryption, + .data_cipher = active_data_cipher_, + .iv_for_data_encryption = iv_for_data_encryption, + .tag_of_data_encryption = {}}; + + return encryption_record; +} + +void storage_core::write_data_to_stream( + util::const_byte_span data, + const optional_binlog_encryption_record &encryption_record, + std::uint64_t offset) { + if (!encryption_record.has_value()) { + // an early return when no encryption is needed + backend_->write_data_to_stream(data); + return; + } + + // as for security reasons our intent is to not store file keys in plaintext + // permanently, we need to decrypt the file key with the KEK before we can + // use it for data encryption. + + const auto &keyring_record{keyring_->get_key(encryption_record->kek_id)}; + + const auto &kek_cipher{keyring_record.get<"cipher">()}; + const auto &kek{keyring_record.get<"data_hex">().get_data()}; + + util::const_byte_span iv_for_file_key_encryption_v{}; + if (encryption_record->iv_for_file_key_encryption.has_value()) { + iv_for_file_key_encryption_v = + *encryption_record->iv_for_file_key_encryption; + }; + util::const_byte_span tag_of_file_key_encryption_v{}; + if (encryption_record->tag_of_file_key_encryption.has_value()) { + tag_of_file_key_encryption_v = + *encryption_record->tag_of_file_key_encryption; + } + + // creating a context for the file key decryption + opensslpp::cipher_context file_key_decryption_context{ + opensslpp::cipher_context_operation_type::decryption, kek_cipher, kek, + iv_for_file_key_encryption_v, tag_of_file_key_encryption_v}; + util::hex_value_storage file_key_decrypted{ + std::size(encryption_record->file_key_encrypted_with_kek)}; + file_key_decryption_context.update( + encryption_record->file_key_encrypted_with_kek, file_key_decrypted); + file_key_decryption_context.finalize(); + + // creating an context for data encryption with the data cipher, the file + // key (decrypted previously), and the IV for data encryption + + auto data_encryption_context{opensslpp::cipher_context::create_with_offset( + offset, opensslpp::cipher_context_operation_type::encryption, + encryption_record->data_cipher, file_key_decrypted, + encryption_record->iv_for_data_encryption)}; + + util::optional_hex_value_storage tag_of_data_encryption{}; + util::byte_span tag_of_data_encryption_v{}; + const auto data_encryption_tag_size{ + data_encryption_context.get_tag_size_in_bytes()}; + if (data_encryption_tag_size != 0U) { + tag_of_data_encryption.emplace(data_encryption_tag_size); + tag_of_data_encryption_v = *tag_of_data_encryption; + } + + util::hex_value_storage encrypted_data{std::size(data)}; + data_encryption_context.update(data, encrypted_data); + data_encryption_context.finalize(tag_of_data_encryption_v); + + backend_->write_data_to_stream(encrypted_data); + + // TODO: update file data encryption tag here, if one day we decide to + // support GCM mode for file data encryption +} + +} // namespace binsrv diff --git a/src/binsrv/storage_core.hpp b/src/binsrv/storage_core.hpp new file mode 100644 index 0000000..5a31bfb --- /dev/null +++ b/src/binsrv/storage_core.hpp @@ -0,0 +1,267 @@ +// Copyright (c) 2023-2024 Percona and/or its affiliates. +// +// This program is free software; you can redistribute it and/or modify +// it under the terms of the GNU General Public License, version 2.0, +// as published by the Free Software Foundation. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU General Public License, version 2.0, for more details. +// +// You should have received a copy of the GNU General Public License +// along with this program; if not, write to the Free Software +// Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA + +#ifndef BINSRV_STORAGE_CORE_HPP +#define BINSRV_STORAGE_CORE_HPP + +#include "binsrv/storage_core_fwd.hpp" // IWYU pragma: export + +#include +#include +#include +#include +#include + +#include "binsrv/basic_keyring_fwd.hpp" +#include "binsrv/basic_logger_fwd.hpp" +#include "binsrv/basic_storage_backend_fwd.hpp" +#include "binsrv/encryption_config_fwd.hpp" +#include "binsrv/encryption_format_type_fwd.hpp" +#include "binsrv/main_config_fwd.hpp" +#include "binsrv/replication_mode_type_fwd.hpp" + +#include "binsrv/events/composite_binlog_name.hpp" + +#include "binsrv/gtids/gtid_fwd.hpp" +#include "binsrv/gtids/gtid_set.hpp" + +#include "binsrv/models/binlog_file_encryption_record_fwd.hpp" + +#include "binsrv/events/common_types.hpp" + +#include "util/byte_span_fwd.hpp" +#include "util/ctime_timestamp_fwd.hpp" +#include "util/ctime_timestamp_range.hpp" +#include "util/hex_value.hpp" + +namespace binsrv { + +struct binlog_encryption_record { + std::string kek_id; + util::hex_value_storage file_key_encrypted_with_kek; + util::optional_hex_value_storage iv_for_file_key_encryption; + util::optional_hex_value_storage tag_of_file_key_encryption; + std::string data_cipher; + util::hex_value_storage iv_for_data_encryption; + util::optional_hex_value_storage tag_of_data_encryption; + + [[nodiscard]] static models::binlog_file_encryption_record + to_model(const binlog_encryption_record &record); + [[nodiscard]] static binlog_encryption_record + from_model(const models::binlog_file_encryption_record &model); +}; + +struct binlog_record { + // binlog file name + events::composite_binlog_name name; + // binlog file size in bytes + std::uint64_t size{0ULL}; + // accumulated GTIDs present in the binlog files before this one + gtids::optional_gtid_set previous_gtids{}; + // GTIDs present in this binlog file + gtids::optional_gtid_set added_gtids{}; + // minimum and maximum event timestamps observed in this binlog file + util::ctime_timestamp_range timestamps{}; + // sequence_number of the last transaction seen in this file - + // used for GTID rewrite-mode resume state persistence + events::seq_no_t last_sequence_number{0ULL}; + // optional encryption parameters + optional_binlog_encryption_record encryption{}; +}; + +class [[nodiscard]] storage_core { +public: + static constexpr std::string_view default_binlog_index_name{"binlog.index"}; + static constexpr std::string_view default_binlog_index_entry_path{"."}; + static constexpr std::string_view metadata_name{"metadata.json"}; + static constexpr std::string_view binlog_metadata_extension{".json"}; + + storage_core(basic_logger_ptr logger, const main_config &config, + storage_construction_mode_type construction_mode); + + storage_core(const storage_core &) = delete; + storage_core &operator=(const storage_core &) = delete; + storage_core(storage_core &&) = delete; + storage_core &operator=(storage_core &&) = delete; + + ~storage_core(); + + [[nodiscard]] storage_construction_mode_type + get_construction_mode() const noexcept { + return construction_mode_; + } + + [[nodiscard]] const gtids::gtid_set &get_purged_gtids() const noexcept { + return purged_gtids_; + } + 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 { + // no need to acquire the mutex as replication_mode_ is immutable after + // construction + return replication_mode_; + } + [[nodiscard]] bool is_in_gtid_replication_mode() const noexcept; + + [[nodiscard]] const binlog_record_container & + get_binlog_records() const noexcept { + return binlog_records_; + } + [[nodiscard]] bool is_empty() const noexcept { + return binlog_records_.empty(); + } + [[nodiscard]] events::composite_binlog_name get_current_binlog_name() const { + return is_empty() ? events::composite_binlog_name{} + : get_current_binlog_record().name; + } + [[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; + } + [[nodiscard]] events::seq_no_t last_sequence_number() const noexcept { + return is_empty() ? 0ULL : get_current_binlog_record().last_sequence_number; + } + + [[nodiscard]] std::uint64_t get_flushed_position() const noexcept { + return is_empty() ? 0ULL : get_current_binlog_record().size; + } + + [[nodiscard]] bool is_binlog_open() const noexcept; + + [[nodiscard]] open_binlog_status + open_binlog(const events::composite_binlog_name &binlog_name); + void write_event_block(util::const_byte_span event_block_data, + const gtids::gtid_set &block_gtids, + const util::ctime_timestamp_range &block_timestamps, + events::seq_no_t block_max_sequence_number); + void close_binlog(); + + // Removes the contiguous prefix of binlog records [front, target] + // (inclusive) from the storage and returns a pair: + // .first - the dropped records in chronological order (oldest + // first), suitable for direct iteration by the caller + // to build a response; + // .second - empty on full success; non-empty when the best-effort + // step-3 cleanup (removal of victim payload + metadata + // objects) failed for at least one object after the + // step-2 index rewrite had already committed. The purge + // itself is considered successful in this case, but the + // storage on disk now contains orphan files that the + // constructor's validators will refuse to open on next + // startup The string carries the underlying cleanup + // error message so the caller. + [[nodiscard]] std::pair + purge_binlogs(const events::composite_binlog_name &target); + + [[nodiscard]] std::string + get_binlog_uri(const events::composite_binlog_name &binlog_name) const; + + [[nodiscard]] bool is_keyring_initialized() const noexcept { + return static_cast(keyring_); + } + [[nodiscard]] std::string get_keyring_description() const; + [[nodiscard]] std::string get_active_kek_description() const; + [[nodiscard]] std::string get_encryption_format_description() const; + + [[nodiscard]] bool has_active_kek() const noexcept { + return !active_kek_id_.empty(); + } + +private: + basic_logger_ptr logger_; + storage_construction_mode_type construction_mode_; + basic_keyring_ptr keyring_; + optional_encryption_format_type encryption_format_; + std::string active_kek_id_; + std::string active_data_cipher_{}; + basic_storage_backend_ptr backend_; + + replication_mode_type replication_mode_; + gtids::gtid_set purged_gtids_{}; + binlog_record_container binlog_records_{}; + + void remove_temporary_objects(storage_object_name_container &object_names); + + void initialize_storage_encryption( + const optional_encryption_config &encryption_config); + + void ensure_streaming_mode() const; + void ensure_purging_mode() const; + + [[nodiscard]] const binlog_record & + get_current_binlog_record() const noexcept { + return binlog_records_.back(); + } + [[nodiscard]] binlog_record &get_current_binlog_record() noexcept { + return binlog_records_.back(); + } + + [[nodiscard]] open_binlog_status open_new_binlog_file_internal( + const events::composite_binlog_name &binlog_name); + [[nodiscard]] open_binlog_status + open_existing_binlog_file_internal(std::uint64_t open_stream_offset); + + void load_binlog_index(); + void validate_binlog_index( + const storage_object_name_container &object_names) const; + void save_binlog_index() const; + + void load_metadata(); + void validate_metadata( + replication_mode_type replication_mode, + const optional_encryption_format_type &encryption_format) const; + void save_metadata() const; + + [[nodiscard]] static std::string generate_binlog_metadata_name( + const events::composite_binlog_name &binlog_name); + [[nodiscard]] binlog_record + load_binlog_metadata(const events::composite_binlog_name &binlog_name) const; + void validate_binlog_metadata(const binlog_record &record) const; + void save_binlog_metadata(const binlog_record &record) const; + + void load_and_validate_binlog_metadata_set( + const storage_object_name_container &object_names, + const storage_object_name_container &object_metadata_names); + + [[nodiscard]] optional_binlog_encryption_record + generate_binlog_encryption_record() const; + + void write_data_to_stream( + util::const_byte_span data, + const optional_binlog_encryption_record &encryption_record, + std::uint64_t offset); +}; + +} // namespace binsrv + +#endif // BINSRV_STORAGE_CORE_HPP diff --git a/src/binsrv/storage_core_fwd.hpp b/src/binsrv/storage_core_fwd.hpp new file mode 100644 index 0000000..2ee7b64 --- /dev/null +++ b/src/binsrv/storage_core_fwd.hpp @@ -0,0 +1,51 @@ +// Copyright (c) 2023-2024 Percona and/or its affiliates. +// +// This program is free software; you can redistribute it and/or modify +// it under the terms of the GNU General Public License, version 2.0, +// as published by the Free Software Foundation. +// +// This program is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU General Public License, version 2.0, for more details. +// +// You should have received a copy of the GNU General Public License +// along with this program; if not, write to the Free Software +// Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA + +#ifndef BINSRV_STORAGE_CORE_FWD_HPP +#define BINSRV_STORAGE_CORE_FWD_HPP + +#include +#include +#include +#include + +namespace binsrv { + +enum class storage_construction_mode_type : std::uint8_t { + querying_only, + streaming, + purging +}; + +enum class open_binlog_status : std::uint8_t { + created, + opened_empty, + opened_at_magic_payload_offset, + opened_with_data_present +}; + +struct binlog_encryption_record; +using optional_binlog_encryption_record = + std::optional; + +struct binlog_record; +using binlog_record_container = std::vector; + +class storage_core; +using storage_core_ptr = std::unique_ptr; + +} // namespace binsrv + +#endif // BINSRV_STORAGE_CORE_FWD_HPP diff --git a/src/binsrv/storage_fwd.hpp b/src/binsrv/storage_fwd.hpp index 8fa29d9..ccb78af 100644 --- a/src/binsrv/storage_fwd.hpp +++ b/src/binsrv/storage_fwd.hpp @@ -16,24 +16,10 @@ #ifndef BINSRV_STORAGE_FWD_HPP #define BINSRV_STORAGE_FWD_HPP -#include #include namespace binsrv { -enum class storage_construction_mode_type : std::uint8_t { - querying_only, - streaming, - purging -}; - -enum class open_binlog_status : std::uint8_t { - created, - opened_empty, - opened_at_magic_payload_offset, - opened_with_data_present -}; - class storage; using storage_ptr = std::shared_ptr; diff --git a/src/operations/collector_context.cpp b/src/operations/collector_context.cpp index 2eb4c69..c1723a3 100644 --- a/src/operations/collector_context.cpp +++ b/src/operations/collector_context.cpp @@ -37,6 +37,7 @@ #include "binsrv/main_config.hpp" #include "binsrv/replication_mode_type.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core_fwd.hpp" #include "binsrv/events/code_type.hpp" #include "binsrv/events/common_header_flag_type.hpp" diff --git a/src/operations/fetch_operation.cpp b/src/operations/fetch_operation.cpp index 82f36e8..f9b3494 100644 --- a/src/operations/fetch_operation.cpp +++ b/src/operations/fetch_operation.cpp @@ -35,6 +35,7 @@ #include "binsrv/log_severity.hpp" #include "binsrv/main_config.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core_fwd.hpp" #include "easymysql/connection_fwd.hpp" diff --git a/src/operations/list_operation.cpp b/src/operations/list_operation.cpp index 322d283..3dd4377 100644 --- a/src/operations/list_operation.cpp +++ b/src/operations/list_operation.cpp @@ -22,6 +22,7 @@ #include "binsrv/main_config.hpp" #include "binsrv/null_logger.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core.hpp" #include "binsrv/models/error_response.hpp" #include "binsrv/models/search_response.hpp" diff --git a/src/operations/model_helpers.cpp b/src/operations/model_helpers.cpp index 04a0727..19e54b8 100644 --- a/src/operations/model_helpers.cpp +++ b/src/operations/model_helpers.cpp @@ -18,15 +18,16 @@ #include #include "binsrv/storage.hpp" +#include "binsrv/storage_core.hpp" #include "binsrv/models/binlog_file_record.hpp" #include "binsrv/models/search_response.hpp" namespace operations { -void append_record_to_search_response( - binsrv::models::search_response &response, const binsrv::storage &storage, - const binsrv::storage::binlog_record &record) { +void append_record_to_search_response(binsrv::models::search_response &response, + const binsrv::storage &storage, + const binsrv::binlog_record &record) { binsrv::models::binlog_file_record record_model{ {{record.name.str()}, {record.size}, @@ -36,8 +37,7 @@ void append_record_to_search_response( {record.timestamps.get_min_timestamp()}, {record.timestamps.get_max_timestamp()}, {record.encryption.has_value() - ? binsrv::storage::binlog_encryption_record::to_model( - *record.encryption) + ? binsrv::binlog_encryption_record::to_model(*record.encryption) : binsrv::models::optional_binlog_file_encryption_record{}}}}; response.add_record(std::move(record_model)); } diff --git a/src/operations/model_helpers.hpp b/src/operations/model_helpers.hpp index cc6b47b..9250a96 100644 --- a/src/operations/model_helpers.hpp +++ b/src/operations/model_helpers.hpp @@ -22,9 +22,9 @@ namespace operations { -void append_record_to_search_response( - binsrv::models::search_response &response, const binsrv::storage &storage, - const binsrv::storage::binlog_record &record); +void append_record_to_search_response(binsrv::models::search_response &response, + const binsrv::storage &storage, + const binsrv::binlog_record &record); } // namespace operations diff --git a/src/operations/pull_operation.cpp b/src/operations/pull_operation.cpp index c28b66f..9b7388c 100644 --- a/src/operations/pull_operation.cpp +++ b/src/operations/pull_operation.cpp @@ -37,6 +37,7 @@ #include "binsrv/log_severity.hpp" #include "binsrv/main_config.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core_fwd.hpp" #include "easymysql/connection_fwd.hpp" diff --git a/src/operations/purge_binlogs_operation.cpp b/src/operations/purge_binlogs_operation.cpp index aefe18f..5807306 100644 --- a/src/operations/purge_binlogs_operation.cpp +++ b/src/operations/purge_binlogs_operation.cpp @@ -22,6 +22,7 @@ #include "binsrv/main_config.hpp" #include "binsrv/null_logger.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core.hpp" #include "binsrv/events/composite_binlog_name.hpp" diff --git a/src/operations/search_by_gtid_set_operation.cpp b/src/operations/search_by_gtid_set_operation.cpp index 6fa2846..ba58144 100644 --- a/src/operations/search_by_gtid_set_operation.cpp +++ b/src/operations/search_by_gtid_set_operation.cpp @@ -23,6 +23,7 @@ #include "binsrv/main_config.hpp" #include "binsrv/null_logger.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core.hpp" #include "binsrv/models/error_response.hpp" #include "binsrv/models/search_response.hpp" diff --git a/src/operations/search_by_timestamp_operation.cpp b/src/operations/search_by_timestamp_operation.cpp index 02829f5..5648a05 100644 --- a/src/operations/search_by_timestamp_operation.cpp +++ b/src/operations/search_by_timestamp_operation.cpp @@ -23,6 +23,7 @@ #include "binsrv/main_config.hpp" #include "binsrv/null_logger.hpp" #include "binsrv/storage.hpp" +#include "binsrv/storage_core.hpp" #include "binsrv/models/error_response.hpp" #include "binsrv/models/search_response.hpp"