Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 8 additions & 2 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
57 changes: 57 additions & 0 deletions src/binsrv/indexed_event_block.cpp
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);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree


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
75 changes: 75 additions & 0 deletions src/binsrv/indexed_event_block.hpp
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
27 changes: 27 additions & 0 deletions src/binsrv/indexed_event_block_fwd.hpp
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
9 changes: 9 additions & 0 deletions src/binsrv/storage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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();
}
Expand Down
7 changes: 7 additions & 0 deletions src/binsrv/storage.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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 {
Expand Down Expand Up @@ -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;
Expand Down
115 changes: 115 additions & 0 deletions src/binsrv/storage_core.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The 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
Expand Down
Loading
Loading