-
Notifications
You must be signed in to change notification settings - Fork 11
PBS-42 feature: Implement basic support for COM_BINLOG_DUMP packet handling (part 3) #194
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
percona-ysorokin
merged 1 commit into
Percona-Lab:main
from
percona-ysorokin:com_binlog_dump_get_event_block
Sep 29, 2026
Merged
Changes from all commits
Commits
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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 <cassert> | ||
| #include <cstddef> | ||
| #include <string> | ||
| #include <utility> | ||
|
|
||
| #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 | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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 <cstddef> | ||
| #include <iterator> | ||
| #include <vector> | ||
|
|
||
| #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<index_record>; | ||
|
|
||
| // 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 |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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 <memory> | ||
| namespace binsrv { | ||
|
|
||
| class indexed_event_block; | ||
| using indexed_event_block_ptr = std::unique_ptr<indexed_event_block>; | ||
|
|
||
| } // namespace binsrv | ||
|
|
||
| #endif // BINSRV_INDEXED_EVENT_BLOCK_FWD_HPP |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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)}; | ||
|
Comment on lines
+532
to
+533
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. As I understand, this is TODO. For now, only unencrypted things work |
||
|
|
||
| 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 | ||
|
|
||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Agree