From f5f5ecea01d81ef54118b77f486f466556fa20d9 Mon Sep 17 00:00:00 2001 From: Yura Sorokin Date: Mon, 28 Sep 2026 17:13:03 +0200 Subject: [PATCH] PBS-42 feature: Implement basic support for COM_BINLOG_DUMP packet handling (part 3) https://perconadev.atlassian.net/browse/PBS-42 Implemented basics for sending binlog events in position-based mode. Current limitations: - Position-based replication only. - Requested ["binlog_name":position] are ignored as if user always requests ["":4] - Blocking mode is not supported - we immediately disconnect when the last available event is sent. - Filtering by 'server_id' is not supported. - Sending heartbeat events is not supported. * Introduced 'binsrv::storage_core::fetch_event_block()' method that tries to read the requested range from the requested binlog file and incapsulates continuation logic. Both reading 'binlog_records' and retrieving data from the backend are done under the shared lock which makes it possible to call this method from multiple network connections in parallel. * Implemented 'binsrv::storage::fetch_event_block()' as a simple redirect to 'binsrv::storage_core::fetch_event_block()'. * Introduced 'binsrv::indexed_event_block' class that identifies event boundaries in the provided memory block and helps to deal with truncated events at the end of the block. * Introduced 'operations::sender_context' class that uses 'binsrv::storage::fetch_event_block()' and 'binsrv::indexed_event_block' behind the scene to generate a continuous sequence of available events. * 'minimysql::network_service::session()' coroutine reworked to send events from the 'operations::sender_context' rather than from the 'minimysql::simple_event_collection'. * Removed 'minimysql::simple_event_collection' class. --- CMakeLists.txt | 10 +- src/binsrv/indexed_event_block.cpp | 57 +++++++ src/binsrv/indexed_event_block.hpp | 75 +++++++++ src/binsrv/indexed_event_block_fwd.hpp | 27 +++ src/binsrv/storage.cpp | 9 + src/binsrv/storage.hpp | 7 + src/binsrv/storage_core.cpp | 115 +++++++++++++ src/binsrv/storage_core.hpp | 56 +++++++ src/minimysql/network_service.cpp | 191 ++++++++++++---------- src/minimysql/sample_event_collection.cpp | 97 ----------- src/minimysql/sample_event_collection.hpp | 43 ----- src/operations/sender_context.cpp | 100 +++++++++++ src/operations/sender_context.hpp | 64 ++++++++ src/operations/sender_context_fwd.hpp | 25 +++ 14 files changed, 646 insertions(+), 230 deletions(-) create mode 100644 src/binsrv/indexed_event_block.cpp create mode 100644 src/binsrv/indexed_event_block.hpp create mode 100644 src/binsrv/indexed_event_block_fwd.hpp delete mode 100644 src/minimysql/sample_event_collection.cpp delete mode 100644 src/minimysql/sample_event_collection.hpp create mode 100644 src/operations/sender_context.cpp create mode 100644 src/operations/sender_context.hpp create mode 100644 src/operations/sender_context_fwd.hpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 42fdc8a..7c5b049 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -525,6 +525,10 @@ set(operations_source_files src/operations/search_by_timestamp_operation.hpp src/operations/search_by_timestamp_operation.cpp + src/operations/sender_context_fwd.hpp + src/operations/sender_context.hpp + src/operations/sender_context.cpp + src/operations/version_operation.hpp src/operations/version_operation.cpp ) @@ -634,6 +638,10 @@ set(binsrv_source_files src/binsrv/s3_storage_backend.hpp src/binsrv/s3_storage_backend.cpp + src/binsrv/indexed_event_block_fwd.hpp + src/binsrv/indexed_event_block.hpp + src/binsrv/indexed_event_block.cpp + src/binsrv/storage_fwd.hpp src/binsrv/storage.hpp src/binsrv/storage.cpp @@ -702,8 +710,6 @@ set(minimysql_source_files src/minimysql/network_io_operations.cpp src/minimysql/network_service.hpp src/minimysql/network_service.cpp - src/minimysql/sample_event_collection.hpp - src/minimysql/sample_event_collection.cpp ) add_executable(binlog_server diff --git a/src/binsrv/indexed_event_block.cpp b/src/binsrv/indexed_event_block.cpp new file mode 100644 index 0000000..908b433 --- /dev/null +++ b/src/binsrv/indexed_event_block.cpp @@ -0,0 +1,57 @@ +// Copyright (c) 2026 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 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/indexed_event_block.hpp" + +#include +#include +#include +#include + +#include "binsrv/events/common_header_view.hpp" + +#include "util/dynamic_byte_buffer_fwd.hpp" + +namespace binsrv { + +indexed_event_block::indexed_event_block(util::dynamic_byte_buffer buffer) + : buffer_(std::move(buffer)) { + auto offset{0UZ}; + while (true) { + const auto remaining_size{std::size(buffer_) - offset}; + if (remaining_size < events::common_header_view_base::size_in_bytes) { + break; + } + + const events::common_header_view header{ + util::const_byte_span{buffer_}.subspan( + offset, events::common_header_view_base::size_in_bytes)}; + const auto event_size{header.get_event_size_raw()}; + + assert(event_size >= events::common_header_view_base::size_in_bytes); + + if (event_size > remaining_size) { + break; + } + + // TODO: consider performing checksum validation here + index_.push_back({offset, event_size}); + offset += event_size; + } + // truncating buffer so that it would only contain complete events + buffer_.resize(offset); +} + +} // namespace binsrv diff --git a/src/binsrv/indexed_event_block.hpp b/src/binsrv/indexed_event_block.hpp new file mode 100644 index 0000000..8e78694 --- /dev/null +++ b/src/binsrv/indexed_event_block.hpp @@ -0,0 +1,75 @@ +// Copyright (c) 2026 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 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_INDEXED_EVENT_BLOCK_HPP +#define BINSRV_INDEXED_EVENT_BLOCK_HPP + +#include "binsrv/indexed_event_block_fwd.hpp" // IWYU pragma: export + +#include +#include +#include + +#include "util/byte_span_fwd.hpp" +#include "util/dynamic_byte_buffer_fwd.hpp" + +namespace binsrv { + +class [[nodiscard]] indexed_event_block { +public: + struct index_record { + std::size_t offset; + std::size_t size; + }; + using index_type = std::vector; + + // Takes ownership of the provided buffer and parses event boundaries within + // it. 'buffer' must start from a valid event position (where common header + // starts) but may end in the middle of an event. After parsing, + // 'get_actual_size()' will return the size of the parsed portion of the + // buffer. + + // deliberately passing by value as we are moving from this argument + explicit indexed_event_block(util::dynamic_byte_buffer buffer); + + indexed_event_block(const indexed_event_block &) = delete; + indexed_event_block(indexed_event_block &&) = delete; + indexed_event_block &operator=(const indexed_event_block &) = delete; + indexed_event_block &operator=(indexed_event_block &&) = delete; + + ~indexed_event_block() = default; + + [[nodiscard]] std::size_t get_actual_size() const noexcept { + return std::size(buffer_); + } + [[nodiscard]] bool is_empty() const noexcept { return std::empty(buffer_); } + + [[nodiscard]] std::size_t get_number_of_events() const noexcept { + return std::size(index_); + } + [[nodiscard]] util::const_byte_span + get_event(std::size_t index) const noexcept { + const auto &record{index_[index]}; + return util::const_byte_span{buffer_}.subspan(record.offset, record.size); + } + +private: + util::dynamic_byte_buffer buffer_; + index_type index_; +}; + +} // namespace binsrv + +#endif // BINSRV_INDEXED_EVENT_BLOCK_HPP diff --git a/src/binsrv/indexed_event_block_fwd.hpp b/src/binsrv/indexed_event_block_fwd.hpp new file mode 100644 index 0000000..ca7a42a --- /dev/null +++ b/src/binsrv/indexed_event_block_fwd.hpp @@ -0,0 +1,27 @@ +// Copyright (c) 2026 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 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_INDEXED_EVENT_BLOCK_FWD_HPP +#define BINSRV_INDEXED_EVENT_BLOCK_FWD_HPP + +#include +namespace binsrv { + +class indexed_event_block; +using indexed_event_block_ptr = std::unique_ptr; + +} // namespace binsrv + +#endif // BINSRV_INDEXED_EVENT_BLOCK_FWD_HPP diff --git a/src/binsrv/storage.cpp b/src/binsrv/storage.cpp index 518c036..2db1460 100644 --- a/src/binsrv/storage.cpp +++ b/src/binsrv/storage.cpp @@ -38,8 +38,10 @@ #include "binsrv/gtids/gtid.hpp" #include "binsrv/gtids/gtid_set.hpp" +#include "util/byte_range_fwd.hpp" #include "util/byte_span.hpp" #include "util/ctime_timestamp.hpp" +#include "util/dynamic_byte_buffer_fwd.hpp" #include "util/exception_location_helpers.hpp" namespace binsrv { @@ -225,6 +227,13 @@ storage::purge_binlogs(const events::composite_binlog_name &target) { return core_->get_binlog_uri(binlog_name); } +[[nodiscard]] bool +storage::fetch_event_block(events::composite_binlog_name &binlog_name, + util::byte_range &range, + util::dynamic_byte_buffer &buffer) const { + return core_->fetch_event_block(binlog_name, range, buffer); +} + [[nodiscard]] std::string storage::get_keyring_description() const { return core_->get_keyring_description(); } diff --git a/src/binsrv/storage.hpp b/src/binsrv/storage.hpp index da6329f..f07c7b4 100644 --- a/src/binsrv/storage.hpp +++ b/src/binsrv/storage.hpp @@ -29,6 +29,7 @@ #include "binsrv/basic_storage_backend_fwd.hpp" #include "binsrv/encryption_config_fwd.hpp" #include "binsrv/encryption_format_type_fwd.hpp" +#include "binsrv/indexed_event_block_fwd.hpp" #include "binsrv/main_config_fwd.hpp" #include "binsrv/replication_mode_type_fwd.hpp" #include "binsrv/storage_core_fwd.hpp" @@ -42,9 +43,11 @@ #include "binsrv/events/common_types.hpp" +#include "util/byte_range_fwd.hpp" #include "util/byte_span_fwd.hpp" #include "util/ctime_timestamp_fwd.hpp" #include "util/ctime_timestamp_range.hpp" +#include "util/dynamic_byte_buffer_fwd.hpp" #include "util/hex_value.hpp" namespace binsrv { @@ -119,6 +122,10 @@ class [[nodiscard]] storage { [[nodiscard]] std::string get_binlog_uri(const events::composite_binlog_name &binlog_name) const; + [[nodiscard]] bool + fetch_event_block(events::composite_binlog_name &binlog_name, + util::byte_range &range, + util::dynamic_byte_buffer &buffer) const; [[nodiscard]] bool is_keyring_initialized() const noexcept; [[nodiscard]] std::string get_keyring_description() const; diff --git a/src/binsrv/storage_core.cpp b/src/binsrv/storage_core.cpp index d9e47bc..b4f559a 100644 --- a/src/binsrv/storage_core.cpp +++ b/src/binsrv/storage_core.cpp @@ -58,8 +58,10 @@ #include "opensslpp/cipher_context.hpp" #include "opensslpp/crypto_rng.hpp" +#include "util/byte_range.hpp" #include "util/byte_span.hpp" #include "util/ctime_timestamp_range.hpp" +#include "util/dynamic_byte_buffer_fwd.hpp" #include "util/exception_location_helpers.hpp" namespace binsrv { @@ -423,6 +425,119 @@ storage_core::purge_binlogs(const events::composite_binlog_name &target) { return {std::move(removed_records), std::move(cleanup_warning_message)}; } +[[nodiscard]] bool +storage_core::fetch_event_block(events::composite_binlog_name &binlog_name, + util::byte_range &range, + util::dynamic_byte_buffer &buffer) const { + static const util::byte_range magic_empty_range{events::magic_binlog_offset, + 0ULL}; + // If the offset in the 'range' is less than + // 'binsrv::events::magic_binlog_offset' (4), the method will return false. + if (range.get_offset() < events::magic_binlog_offset) { + return false; + } + + // If the specified 'range' is an open range (has no length set), this + // method will return false. + if (!range.has_length()) { + return false; + } + + // If the specified 'binlog_name' is an empty object and offset of the + // 'range' is not equal to 'binsrv::events::magic_binlog_offset' (4), + // the method will return false. + if (binlog_name.is_empty() && + range.get_offset() != events::magic_binlog_offset) { + return false; + } + + const std::shared_lock lock{mutex_}; + + // If the specified 'binlog_name' is an empty object and offset of the + // 'range' is equal to 'binsrv::events::magic_binlog_offset' (4), and + // storage has no binlog records, the method will return true, + // will set binlog name to an empty object, range to "[4; 0]", + // and buffer to an empty buffer. + if (binlog_records_.empty()) { + // EOF is returned only when 'binlog_name' is an empty object + if (!binlog_name.is_empty()) { + return false; + } + range = magic_empty_range; + buffer.clear(); + return true; + } + + binlog_record_container::const_iterator record_it{}; + // If the specified 'binlog_name' is an empty object and offset of the + // 'range' is equal to 'binsrv::events::magic_binlog_offset' (4), and + // there is at least one binlog record available, when checking other + // rules, we will assume that 'binlog_name' from now on will be equal to + // the first available binlog file name. + if (binlog_name.is_empty()) { + record_it = std::cbegin(binlog_records_); + } else { + // If the specified (or resolved) 'binlog_name' does not exist in storage, + // the method will return false. + record_it = std::ranges::find(std::as_const(binlog_records_), binlog_name, + &binlog_record::name); + if (record_it == std::cend(binlog_records_)) { + return false; + } + } + auto resolved_binlog_name{record_it->name}; + + // If 'range' is an empty range, the method will return true without + // attempting to read any data. The range will remain unchanged, the + // buffer will be set to an empty object, and 'binlog_name' will be changed + // only if it was originally empty and was resolved to the first available + // binlog file. + if (range.is_empty()) { + binlog_name = std::move(resolved_binlog_name); + buffer.clear(); + return true; + } + + // If the specified 'range' has an offset that is beyond the end of the + // specified binlog, the method will return false. + auto read_offset{range.get_offset()}; + if (read_offset > record_it->size) { + return false; + } + + // If the 'range.get_offset()' is equal to the length of the binlog file + // specified by the 'binlog_name', this method will return true and + // will try to read 'range.get_length()' bytes from the + // offset 'binsrv::events::magic_binlog_offset' (4) of the next binlog + // file, if available. 'range' and 'binlog_name' will be updated + // accordingly. + if (read_offset == record_it->size) { + ++record_it; + // If the next file is not available, the method will return true and will + // leave 'binlog_name' as is, change the 'length' component of the 'range' + // to 0, and set 'buffer' to an empty buffer, indicating EOF. + if (record_it == std::cend(binlog_records_)) { + range = util::byte_range{read_offset, 0ULL}; + buffer.clear(); + return true; + } + resolved_binlog_name = record_it->name; + read_offset = events::magic_binlog_offset; + } + + const std::uint64_t read_length{ + std::min(range.get_length(), record_it->size - read_offset)}; + const util::byte_range resolved_range{read_offset, read_length}; + + auto result_buffer{ + backend_->get_object(resolved_binlog_name.str(), resolved_range)}; + + binlog_name = std::move(resolved_binlog_name); + range = resolved_range; + buffer = std::move(result_buffer); + return true; +} + [[nodiscard]] std::string storage_core::get_binlog_uri( const events::composite_binlog_name &binlog_name) const { // no mutex protection needed as this method calls a const diff --git a/src/binsrv/storage_core.hpp b/src/binsrv/storage_core.hpp index 52df796..2f2650a 100644 --- a/src/binsrv/storage_core.hpp +++ b/src/binsrv/storage_core.hpp @@ -31,6 +31,7 @@ #include "binsrv/basic_storage_backend_fwd.hpp" #include "binsrv/encryption_config_fwd.hpp" #include "binsrv/encryption_format_type_fwd.hpp" +#include "binsrv/indexed_event_block_fwd.hpp" #include "binsrv/main_config_fwd.hpp" #include "binsrv/replication_mode_type_fwd.hpp" @@ -43,9 +44,11 @@ #include "binsrv/events/common_types.hpp" +#include "util/byte_range_fwd.hpp" #include "util/byte_span_fwd.hpp" #include "util/ctime_timestamp_fwd.hpp" #include "util/ctime_timestamp_range.hpp" +#include "util/dynamic_byte_buffer_fwd.hpp" #include "util/hex_value.hpp" namespace binsrv { @@ -179,6 +182,59 @@ class [[nodiscard]] storage_core { [[nodiscard]] std::pair purge_binlogs(const events::composite_binlog_name &target); + // This method will try to read a block of events of length + // 'range.get_length()' from the specified binlog file 'binlog_name', + // starting from 'range.get_offset()'. + // Returns true if the operation was successful, false otherwise. + // Both 'binlog_name' and 'range' are inout parameters and they will + // be updated to reflect the actual portion of the binlog that was read + // if the operation was successful. + // All parameters will remain untouched if the operation fails. + // Special cases: + // - If the offset in the 'range' is less than + // 'binsrv::events::magic_binlog_offset' (4), the method will return false. + // - If the specified 'range' is an open range (has no length set), this + // method will return false. + // - If the specified 'binlog_name' is an empty object and offset of the + // 'range' is not equal to 'binsrv::events::magic_binlog_offset' (4), + // the method will return false. + // - If the specified 'binlog_name' is an empty object and offset of the + // 'range' is equal to 'binsrv::events::magic_binlog_offset' (4), and + // storage has no binlog records, the method will return true, + // will set binlog name to an empty object, range to "[4; 0]", + // and buffer to an empty buffer. + // - If the specified 'binlog_name' is an empty object and offset of the + // 'range' is equal to 'binsrv::events::magic_binlog_offset' (4), and + // there is at least one binlog record available, when checking other + // rules, we will assume that 'binlog_name' from now on will be equal to + // the first available binlog file name. + // - If the specified (or resolved) 'binlog_name' does not exist in storage, + // the method will return false. + // - If 'range' is an empty range, the method will return true without + // attempting to read any data. The range will remain unchanged, the + // buffer will be set to an empty object, and 'binlog_name' will be changed + // only if it was originally empty and was resolved to the first available + // binlog file. + // - If the specified 'range' has an offset that is beyond the end of the + // specified binlog, the method will return false. + // - If the specified 'range' has valid offset for the given 'binlog_name', + // but the length extends beyond the end of the binlog, the method will + // return true and will read only the available portion and update + // 'range' to reflect the actual portion read. + // - If the 'range.get_offset()' is equal to the length of the binlog file + // specified by the 'binlog_name', this method will return true and + // will try to read 'range.get_length()' bytes from the + // offset 'binsrv::events::magic_binlog_offset' (4) of the next binlog + // file, if available. 'range' and 'binlog_name' will be updated + // accordingly. + // If the next file is not available, the method will return true and will + // leave 'binlog_name' as is, change the 'length' component of the 'range' + // to 0, and set 'buffer' to an empty buffer, indicating EOF. + [[nodiscard]] bool + fetch_event_block(events::composite_binlog_name &binlog_name, + util::byte_range &range, + util::dynamic_byte_buffer &buffer) const; + [[nodiscard]] std::string get_binlog_uri(const events::composite_binlog_name &binlog_name) const; diff --git a/src/minimysql/network_service.cpp b/src/minimysql/network_service.cpp index e2984ab..30a566a 100644 --- a/src/minimysql/network_service.cpp +++ b/src/minimysql/network_service.cpp @@ -68,7 +68,8 @@ #include "minimysql/connection_context.hpp" #include "minimysql/network_io_operations.hpp" -#include "minimysql/sample_event_collection.hpp" + +#include "operations/sender_context.hpp" #include "util/byte_span.hpp" @@ -280,9 +281,9 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // parses client greeting [[nodiscard]] boost::asio::awaitable session( // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) - binsrv::basic_logger &logger, + const binsrv::basic_logger_ptr &logger, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) - binsrv::storage &storage, boost::asio::ip::tcp::socket socket, + const binsrv::storage_ptr &storage, boost::asio::ip::tcp::socket socket, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) const binsrv::replication_source_config &cfg) { boost::system::error_code session_ec; @@ -290,7 +291,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { const auto remote_endpoint_str{ boost::lexical_cast(remote_endpoint)}; - const scope_tracer tracer(logger, "session " + remote_endpoint_str); + const scope_tracer tracer(*logger, "session " + remote_endpoint_str); const std::chrono::seconds read_timeout{cfg.get<"read_timeout">()}; const std::chrono::seconds write_timeout{cfg.get<"write_timeout">()}; @@ -314,12 +315,12 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // "caching_sha2_password" const auto server_greeting{context.generate_encoded_server_greeting()}; - print_server_greeting(logger, remote_endpoint, context); + print_server_greeting(*logger, remote_endpoint, context); co_await minimysql::async_write_mysql_frame(socket, server_greeting, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server greeting ({} bytes to {})", - std::size(server_greeting), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : sent server greeting ({} bytes to {})", + std::size(server_greeting), remote_endpoint_str); // receiving and parsing client greeting packet: // capabilities @@ -332,84 +333,85 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // attributes co_await minimysql::async_read_mysql_frame(socket, data, read_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : received client greeting ({} bytes from {})", - std::size(data), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : received client greeting ({} bytes from {})", + std::size(data), remote_endpoint_str); context.parse_client_greeting(data); - print_client_greeting(logger, remote_endpoint, context); + print_client_greeting(*logger, remote_endpoint, context); if (!context.check_shared_plugin_auth_supported()) { - logger.log(binsrv::log_severity::warning, - "net : client does not support plugin authentication"); + logger->log(binsrv::log_severity::warning, + "net : client does not support plugin authentication"); const auto access_denied{context.generate_encoded_access_denied()}; - print_error(logger, remote_endpoint, context, "plugin auth required"); + print_error(*logger, remote_endpoint, context, "plugin auth required"); co_await minimysql::async_write_mysql_frame(socket, access_denied, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server access denied ({} bytes to {})", - std::size(access_denied), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : sent server access denied ({} bytes to {})", + std::size(access_denied), remote_endpoint_str); co_return; } if (context.get_client_auth_method() != context.get_server_auth_method()) { - logger.log_format(binsrv::log_severity::info, - "net : client requested {} authentication that does " - "not match the one " - "associated with the user account ({})", - context.get_client_auth_method(), - context.get_server_auth_method()); + logger->log_format( + binsrv::log_severity::info, + "net : client requested {} authentication that does " + "not match the one " + "associated with the user account ({})", + context.get_client_auth_method(), context.get_server_auth_method()); const auto auth_method_switch{ context.generate_encoded_auth_method_switch()}; - print_generic(logger, remote_endpoint, context, "auth method switch"); + print_generic(*logger, remote_endpoint, context, "auth method switch"); co_await minimysql::async_write_mysql_frame(socket, auth_method_switch, write_timeout); - logger.log_format( + logger->log_format( binsrv::log_severity::debug, "net : sent server auth method switch ({} bytes to {})", std::size(auth_method_switch), remote_endpoint_str); co_await minimysql::async_read_mysql_frame(socket, data, read_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : received client auth method switch response " - "({} bytes from {})", - std::size(data), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : received client auth method switch response " + "({} bytes from {})", + std::size(data), remote_endpoint_str); context.parse_client_auth_method_switch(data); - print_client_auth_method_switch(logger, remote_endpoint, context); + print_client_auth_method_switch(*logger, remote_endpoint, context); } if (!context.check_client_authentication()) { - logger.log_format(binsrv::log_severity::warning, - "net : client authentication failed for {}", - context.get_client_username()); + logger->log_format(binsrv::log_severity::warning, + "net : client authentication failed for {}", + context.get_client_username()); const auto access_denied{context.generate_encoded_access_denied()}; - print_error(logger, remote_endpoint, context, "auth failure"); + print_error(*logger, remote_endpoint, context, "auth failure"); co_await minimysql::async_write_mysql_frame(socket, access_denied, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server access denied ({} bytes to {})", - std::size(access_denied), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : sent server access denied ({} bytes to {})", + std::size(access_denied), remote_endpoint_str); co_return; } - logger.log_format(binsrv::log_severity::info, - "net : client authentication succeeded for {}", - context.get_client_username()); + logger->log_format(binsrv::log_severity::info, + "net : client authentication succeeded for {}", + context.get_client_username()); // sending fast auth success const auto fast_auth_success{context.generate_encoded_fast_auth()}; - print_generic(logger, remote_endpoint, context, + print_generic(*logger, remote_endpoint, context, "auth method data (fast auth)"); co_await minimysql::async_write_mysql_frame(socket, fast_auth_success, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server fast auth success ({} bytes to {})", - std::size(fast_auth_success), remote_endpoint_str); + logger->log_format( + binsrv::log_severity::debug, + "net : sent server fast auth success ({} bytes to {})", + std::size(fast_auth_success), remote_endpoint_str); // sending server ok after successful authentication const auto auth_ok{context.generate_encoded_ok()}; - print_generic(logger, remote_endpoint, context, "ok (auth)"); + print_generic(*logger, remote_endpoint, context, "ok (auth)"); co_await minimysql::async_write_mysql_frame(socket, auth_ok, write_timeout); - logger.log_format( + logger->log_format( binsrv::log_severity::debug, "net : sent server ok after authentication ({} bytes to {})", std::size(auth_ok), remote_endpoint_str); @@ -474,11 +476,11 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { while (!terminated) { context.enter_command_loop_iteration(); co_await minimysql::async_read_mysql_frame(socket, data, read_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : received client command ({} bytes from {})", - std::size(data), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : received client command ({} bytes from {})", + std::size(data), remote_endpoint_str); context.parse_client_command(data); - print_client_command(logger, remote_endpoint, context); + print_client_command(*logger, remote_endpoint, context); switch (context.get_client_mysql_command()) { case minimysql::client_command_type::query: { @@ -486,19 +488,19 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { known_queries.find(context.get_client_statement())}; if (known_query_it != std::end(known_queries)) { const auto resultset{known_query_it->second(context)}; - print_generic(logger, remote_endpoint, context, "resultset"); + print_generic(*logger, remote_endpoint, context, "resultset"); co_await minimysql::async_write_mysql_frames(socket, resultset, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server resultset ({} frames to {})", - std::size(resultset), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : sent server resultset ({} frames to {})", + std::size(resultset), remote_endpoint_str); } else { // return 'syntax error' for every other query const auto syntax_error = context.generate_encoded_syntax_error(); - print_error(logger, remote_endpoint, context, "syntax error"); + print_error(*logger, remote_endpoint, context, "syntax error"); co_await minimysql::async_write_mysql_frame(socket, syntax_error, write_timeout); - logger.log_format( + logger->log_format( binsrv::log_severity::debug, "net : sent server syntax error ({} bytes to {})", std::size(syntax_error), remote_endpoint_str); @@ -506,34 +508,48 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { } break; case minimysql::client_command_type::ping: { const auto ok_after_ping{context.generate_encoded_ok()}; - print_generic(logger, remote_endpoint, context, "ok (ping success)"); + print_generic(*logger, remote_endpoint, context, "ok (ping success)"); co_await minimysql::async_write_mysql_frame(socket, ok_after_ping, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server ok after ping ({} bytes to {})", - std::size(ok_after_ping), remote_endpoint_str); + logger->log_format( + binsrv::log_severity::debug, + "net : sent server ok after ping ({} bytes to {})", + std::size(ok_after_ping), remote_endpoint_str); } break; case minimysql::client_command_type::binlog_dump: { - const minimysql::sample_event_collection sample_events; - // TODO: rework with reading real data from storage; - (void)storage; - for (const auto &event_data : sample_events.get_events()) { - const auto event{context.generate_encoded_binlog_event( - util::as_const_byte_span(event_data))}; - print_generic(logger, remote_endpoint, context, "binlog event"); - co_await minimysql::async_write_mysql_frame(socket, event, + static constexpr auto block_size{1048576UZ}; + + // TODO: initialize sender_context with binlog_name:position extracted + // from the COM_BINLOG_DUMP command. + operations::sender_context sender_ctx{logger, storage, block_size}; + bool fetch_result{}; + util::const_byte_span event_data{}; + while ((fetch_result = sender_ctx.get_event(event_data)) && + !event_data.empty()) { + const auto event_frame{ + context.generate_encoded_binlog_event(event_data)}; + print_generic(*logger, remote_endpoint, context, "binlog event"); + co_await minimysql::async_write_mysql_frame(socket, event_frame, write_timeout); - logger.log_format( + logger->log_format( binsrv::log_severity::debug, "net : sent server binlog event ({} bytes to {})", - std::size(event), remote_endpoint_str); + std::size(event_frame), remote_endpoint_str); } + if (!fetch_result) { + logger->log_format(binsrv::log_severity::error, + "net : failed to fetch next event block for {}", + remote_endpoint_str); + terminated = true; + break; + } + const auto eof = context.generate_encoded_eof(); - print_generic(logger, remote_endpoint, context, "binlog eof"); + print_generic(*logger, remote_endpoint, context, "binlog eof"); co_await minimysql::async_write_mysql_frame(socket, eof, write_timeout); - logger.log_format(binsrv::log_severity::debug, - "net : sent server eof ({} bytes to {})", - std::size(eof), remote_endpoint_str); + logger->log_format(binsrv::log_severity::debug, + "net : sent server eof ({} bytes to {})", + std::size(eof), remote_endpoint_str); terminated = true; } break; case minimysql::client_command_type::quit: { @@ -544,10 +560,10 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { default: { const auto unknown_command_error = context.generate_encoded_unknown_command(); - print_error(logger, remote_endpoint, context, "unknown command"); + print_error(*logger, remote_endpoint, context, "unknown command"); co_await minimysql::async_write_mysql_frame( socket, unknown_command_error, write_timeout); - logger.log_format( + logger->log_format( binsrv::log_severity::debug, "net : sent server unknown command ({} bytes to {})", std::size(unknown_command_error), remote_endpoint_str); @@ -555,7 +571,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { } } } catch (...) { - handle_exception(logger, "session " + remote_endpoint_str); + handle_exception(*logger, "session " + remote_endpoint_str); } } #pragma GCC diagnostic pop @@ -564,14 +580,14 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // coroutine for each accepted connection [[nodiscard]] boost::asio::awaitable listener( // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) - binsrv::basic_logger &logger, + const binsrv::basic_logger_ptr &logger, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) - binsrv::storage &storage, + const binsrv::storage_ptr &storage, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) boost::asio::ip::tcp::acceptor &acceptor, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) const binsrv::replication_source_config &cfg) { - const scope_tracer tracer(logger, "listener"); + const scope_tracer tracer(*logger, "listener"); auto executor = acceptor.get_executor(); @@ -585,14 +601,14 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { listener_ec != boost::asio::error::bad_descriptor) { throw boost::system::system_error{listener_ec}; } - logger.log(binsrv::log_severity::info, "net : listener stopped"); + logger->log(binsrv::log_severity::info, "net : listener stopped"); break; } const auto remote_endpoint{socket.remote_endpoint(listener_ec)}; - logger.log_format(binsrv::log_severity::info, - "net : accepted connection from {}", - boost::lexical_cast(remote_endpoint)); + logger->log_format(binsrv::log_severity::info, + "net : accepted connection from {}", + boost::lexical_cast(remote_endpoint)); // NOLINTNEXTLINE(misc-include-cleaner) boost::asio::co_spawn(executor, @@ -600,7 +616,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { boost::asio::detached); } } catch (...) { - handle_exception(logger, "listener"); + handle_exception(*logger, "listener"); } } @@ -617,8 +633,7 @@ network_service::network_service(binsrv::basic_logger_ptr logger, cfg.get<"port">()})} { assert(logger_); // NOLINTNEXTLINE(misc-include-cleaner) - boost::asio::co_spawn(*context_, - listener(*logger_, *storage_, *acceptor_, cfg), + boost::asio::co_spawn(*context_, listener(logger_, storage_, *acceptor_, cfg), boost::asio::detached); } diff --git a/src/minimysql/sample_event_collection.cpp b/src/minimysql/sample_event_collection.cpp deleted file mode 100644 index 0ada88c..0000000 --- a/src/minimysql/sample_event_collection.cpp +++ /dev/null @@ -1,97 +0,0 @@ -// Copyright (c) 2023-2026 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 "minimysql/sample_event_collection.hpp" - -#include -#include - -#include - -namespace minimysql { - -namespace { - -constexpr std::string_view rotate_artificial{ - "00 00 00 00 04 01 00 00 00 28 00 00 00 00 00 00" - "00 20 00 04 00 00 00 00 00 00 00 62 69 6e 6c 6f" - "67 2e 30 30 30 30 30 31 "}; -constexpr std::string_view fde{ - "1a b2 1c 6a 0f 01 00 00 00 7b 00 00 00 7f 00 00" - "00 00 00 04 00 38 2e 34 2e 37 2d 37 00 00 00 00" - "00 00 00 00 00 00 00 00 00 00 00 00 00 00 00 00" - "00 00 00 00 00 00 00 00 00 00 00 00 00 00 00 00" - "00 00 00 00 00 00 00 1a b2 1c 6a 13 00 0d 00 08" - "00 00 00 00 04 00 04 00 00 00 63 00 04 1a 08 00" - "00 00 00 00 00 02 00 00 00 0a 0a 0a 2a 2a 00 12" - "34 00 0a 28 00 00 01 ee c6 67 82 "}; -constexpr std::string_view previous_gtids{ - "1a b2 1c 6a 23 01 00 00 00 1f 00 00 00 9e 00 00" - "00 80 00 00 00 00 00 00 00 00 00 09 5d aa 16 "}; -constexpr std::string_view first_gtid_log{ - "52 b2 1c 6a 21 01 00 00 00 4d 00 00 00 eb 00 00" - "00 00 00 01 b1 14 55 0c 5d 3d 11 f1 a7 3a 00 15" - "5d f1 f4 92 01 00 00 00 00 00 00 00 02 00 00 00" - "00 00 00 00 00 01 00 00 00 00 00 00 00 9c 7e f6" - "5f 24 53 06 cb 17 3a 01 00 80 dc f9 5c "}; -constexpr std::string_view first_query{ - "52 b2 1c 6a 02 01 00 00 00 7e 00 00 00 69 01 00" - "00 00 00 09 00 00 00 00 00 00 00 04 00 00 2f 00" - "00 00 00 00 00 01 20 00 a0 45 00 00 00 00 06 03" - "73 74 64 04 ff 00 ff 00 ff 00 0c 01 74 65 73 74" - "00 11 06 00 00 00 00 00 00 00 12 ff 00 13 00 74" - "65 73 74 00 43 52 45 41 54 45 20 54 41 42 4c 45" - "20 74 31 28 69 64 20 53 45 52 49 41 4c 20 50 52" - "49 4d 41 52 59 20 4b 45 59 29 d9 1d 2b 2b "}; -constexpr std::string_view second_gtid_log{ - "57 b2 1c 6a 21 01 00 00 00 4f 00 00 00 b8 01 00" - "00 00 00 00 b1 14 55 0c 5d 3d 11 f1 a7 3a 00 15" - "5d f1 f4 92 02 00 00 00 00 00 00 00 02 01 00 00" - "00 00 00 00 00 02 00 00 00 00 00 00 00 8a 35 4f" - "60 24 53 06 fc 15 01 17 3a 01 00 24 35 52 71 "}; -constexpr std::string_view second_query{ - "57 b2 1c 6a 02 01 00 00 00 4b 00 00 00 03 02 00" - "00 08 00 09 00 00 00 00 00 00 00 04 00 00 1d 00" - "00 00 00 00 00 01 20 00 a0 45 00 00 00 00 06 03" - "73 74 64 04 ff 00 ff 00 ff 00 12 ff 00 74 65 73" - "74 00 42 45 47 49 4e c8 4c ce 80 "}; -constexpr std::string_view second_table_map{ - "57 b2 1c 6a 13 01 00 00 00 30 00 00 00 33 02 00" - "00 00 00 73 00 00 00 00 00 01 00 04 74 65 73 74" - "00 02 74 31 00 01 08 00 00 01 01 80 21 d6 bc 8c"}; -constexpr std::string_view second_write_rows{ - "57 b2 1c 6a 1e 01 00 00 00 2c 00 00 00 5f 02 00" - "00 00 00 73 00 00 00 00 00 01 00 02 00 01 ff 00" - "01 00 00 00 00 00 00 00 82 65 a7 86 "}; -constexpr std::string_view second_xid{ - "57 b2 1c 6a 10 01 00 00 00 1f 00 00 00 7e 02 00" - "00 00 00 07 00 00 00 00 00 00 00 e3 36 06 8f "}; - -std::string hex_to_bin(std::string_view hex) { - std::string filtered_hex{hex}; - std::erase(filtered_hex, ' '); - return boost::algorithm::unhex(filtered_hex); -} - -} // anonymous namespace - -sample_event_collection::sample_event_collection() - : events_{hex_to_bin(rotate_artificial), hex_to_bin(fde), - hex_to_bin(previous_gtids), hex_to_bin(first_gtid_log), - hex_to_bin(first_query), hex_to_bin(second_gtid_log), - hex_to_bin(second_query), hex_to_bin(second_table_map), - hex_to_bin(second_write_rows), hex_to_bin(second_xid)} {} - -} // namespace minimysql diff --git a/src/minimysql/sample_event_collection.hpp b/src/minimysql/sample_event_collection.hpp deleted file mode 100644 index 4d9ff1e..0000000 --- a/src/minimysql/sample_event_collection.hpp +++ /dev/null @@ -1,43 +0,0 @@ -// Copyright (c) 2023-2026 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 MINIMYSQL_SAMPLE_EVENT_COLLECTION_HPP -#define MINIMYSQL_SAMPLE_EVENT_COLLECTION_HPP - -#include -#include - -namespace minimysql { - -class sample_event_collection { -public: - using event_data_type = std::string; - static constexpr auto number_of_predefined_events{10UZ}; - using event_data_container = - std::array; - - sample_event_collection(); - - [[nodiscard]] const event_data_container &get_events() const noexcept { - return events_; - } - -private: - event_data_container events_; -}; - -} // namespace minimysql - -#endif // MINIMYSQL_SAMPLE_EVENT_COLLECTION_HPP diff --git a/src/operations/sender_context.cpp b/src/operations/sender_context.cpp new file mode 100644 index 0000000..0050765 --- /dev/null +++ b/src/operations/sender_context.cpp @@ -0,0 +1,100 @@ +// 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 "operations/sender_context.hpp" + +#include +#include +#include +#include +#include + +#include "binsrv/basic_logger.hpp" +#include "binsrv/indexed_event_block.hpp" +#include "binsrv/log_severity.hpp" +#include "binsrv/storage.hpp" + +#include "binsrv/events/protocol_traits_fwd.hpp" + +#include "util/byte_span_fwd.hpp" +#include "util/dynamic_byte_buffer_fwd.hpp" + +namespace operations { + +sender_context::sender_context(binsrv::basic_logger_ptr logger, + binsrv::storage_ptr storage, + std::size_t block_size) + : logger_{std::move(logger)}, storage_{std::move(storage)}, + block_size_{block_size}, + range_{binsrv::events::magic_binlog_offset, block_size_} { + assert(block_size_ > 0UZ); + assert(storage_); + assert(logger_); +} + +sender_context::~sender_context() = default; + +[[nodiscard]] bool sender_context::get_event(util::const_byte_span &event) { + // if this is the very first call when 'event_block_' is not yet set or we + // have consumed all events in the current block + if (!event_block_ || event_index_ == event_block_->get_number_of_events()) { + // early reset to free memory from the previous event block + event_block_.reset(); + + util::dynamic_byte_buffer buffer{}; + // on success both 'binlog_name_' and 'range_' will be updated + if (!storage_->fetch_event_block(binlog_name_, range_, buffer)) { + return false; + } + if (buffer.empty()) { + logger_->log(binsrv::log_severity::info, "sender : fetched EOF"); + // resetting the sender context to its initial state on EOF + binlog_name_ = binsrv::events::composite_binlog_name{}; + range_ = + util::byte_range{binsrv::events::magic_binlog_offset, block_size_}; + event_index_ = 0UZ; + + // setting the event span to an empty object to indicate EOF + event = util::const_byte_span{}; + return true; // EOF + } + logger_->log_format(binsrv::log_severity::info, + "sender : fetched event block of size {}, {}:{}", + std::size(buffer), binlog_name_.str(), + range_.to_string()); + event_block_ = + std::make_unique(std::move(buffer)); + if (event_block_->is_empty()) { + // in case when the received buffer is not empty but after parsing it has + // no valid events, we should treat this situation as an error + + // TODO: double the length of the requested block and retry fetching + return false; + } + event_index_ = 0UZ; + logger_->log_format( + binsrv::log_severity::info, + "sender : parsed event block with {} events, actual size {} byte(s)", + event_block_->get_number_of_events(), event_block_->get_actual_size()); + // preparing 'range_' for the next fetch + range_ = util::byte_range{ + range_.get_offset() + event_block_->get_actual_size(), block_size_}; + } + event = event_block_->get_event(event_index_); + ++event_index_; + return true; +} + +} // namespace operations diff --git a/src/operations/sender_context.hpp b/src/operations/sender_context.hpp new file mode 100644 index 0000000..8a105b9 --- /dev/null +++ b/src/operations/sender_context.hpp @@ -0,0 +1,64 @@ +// 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 OPERATIONS_SENDER_CONTEXT_HPP +#define OPERATIONS_SENDER_CONTEXT_HPP + +#include "operations/sender_context_fwd.hpp" // IWYU pragma: export + +#include + +#include "binsrv/basic_logger_fwd.hpp" +#include "binsrv/indexed_event_block_fwd.hpp" +#include "binsrv/storage_fwd.hpp" + +#include "binsrv/events/composite_binlog_name.hpp" + +#include "util/byte_range.hpp" +#include "util/byte_span_fwd.hpp" + +namespace operations { + +class sender_context { +public: + // deliberately passing by value as we will be moving from these objects + sender_context(binsrv::basic_logger_ptr logger, binsrv::storage_ptr storage, + std::size_t block_size); + + sender_context(const sender_context &) = delete; + sender_context &operator=(const sender_context &) = delete; + sender_context(sender_context &&) = delete; + sender_context &operator=(sender_context &&) = delete; + ~sender_context(); + + // returns false on error + // returns true and sets the event span to a non-empty value on success + // returns true and sets the event span to an empty object on EOF + [[nodiscard]] bool get_event(util::const_byte_span &event); + +private: + binsrv::basic_logger_ptr logger_{}; + binsrv::storage_ptr storage_{}; + + std::size_t block_size_{}; + binsrv::events::composite_binlog_name binlog_name_{}; + util::byte_range range_{}; + binsrv::indexed_event_block_ptr event_block_{}; + std::size_t event_index_{}; +}; + +} // namespace operations + +#endif // OPERATIONS_SENDER_CONTEXT_HPP diff --git a/src/operations/sender_context_fwd.hpp b/src/operations/sender_context_fwd.hpp new file mode 100644 index 0000000..0fe542c --- /dev/null +++ b/src/operations/sender_context_fwd.hpp @@ -0,0 +1,25 @@ +// 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 OPERATIONS_SENDER_CONTEXT_FWD_HPP +#define OPERATIONS_SENDER_CONTEXT_FWD_HPP + +namespace operations { + +class sender_context; + +} // namespace operations + +#endif // OPERATIONS_SENDER_CONTEXT_FWD_HPP