PBS-42 feature: Implement basic support for COM_BINLOG_DUMP packet handling (part 3) - #194
Conversation
…ndling (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.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Encrypted storage, oversized events, malformed headers, and blocking storage reads are not handled safely.
Review effort: Balanced
Findings: 2
Open (5)
Malformed zero-size event header can cause infinite indexing loop · New Encrypted binlog ranges are not decrypted before event parsing · New Shared lock held during backend I/O blocks binlog ingestion · New Synchronous storage reads block the single-threaded application · New Fixed 1 MiB fetch rejects valid larger binlog events · New
What changed in this PR
Adds storage-backed, position-based COM_BINLOG_DUMP event streaming.
Changes:
- Adds block fetching and cross-binlog continuation.
- Adds event-boundary indexing and sender state management.
- Replaces sample events with stored binlog events.
| File | Description |
|---|---|
CMakeLists.txt |
Updates source registration. |
src/operations/sender_context_fwd.hpp |
Declares sender context. |
src/operations/sender_context.hpp |
Defines sender state and API. |
src/operations/sender_context.cpp |
Implements block-to-event iteration. |
src/minimysql/sample_event_collection.hpp |
Removes sample event API. |
src/minimysql/sample_event_collection.cpp |
Removes hard-coded events. |
src/minimysql/network_service.cpp |
Streams storage-backed events. |
src/binsrv/storage.hpp |
Exposes block fetching. |
src/binsrv/storage.cpp |
Delegates block fetching. |
src/binsrv/storage_core.hpp |
Documents the fetch contract. |
src/binsrv/storage_core.cpp |
Implements range and file continuation. |
src/binsrv/indexed_event_block_fwd.hpp |
Declares indexed blocks. |
src/binsrv/indexed_event_block.hpp |
Defines event indexing. |
src/binsrv/indexed_event_block.cpp |
Parses event boundaries. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| 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); |
| auto result_buffer{ | ||
| backend_->get_object(resolved_binlog_name.str(), resolved_range)}; |
There was a problem hiding this comment.
As I understand, this is TODO. For now, only unencrypted things work
| return false; | ||
| } | ||
|
|
||
| const std::shared_lock lock{mutex_}; |
| while ((fetch_result = sender_ctx.get_event(event_data)) && | ||
| !event_data.empty()) { |
| 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; |
There was a problem hiding this comment.
I'm not sure if the proposed "retry with double buffer, then double again if still not enough" is a good strategy. We already know the worst case, why not to just handle it?
I think we should implement either way, rather than being OK, with having 1MiB limit or error
|
|
||
| util::dynamic_byte_buffer buffer{}; | ||
| // on success both 'binlog_name_' and 'range_' will be updated | ||
| if (!storage_->fetch_event_block(binlog_name_, range_, buffer)) { |
There was a problem hiding this comment.
Let's consider the case where small event is followed by big event. Only small event fits the indexed buffer, and all following data is discarded. Once the small event is sent, the loop re-fetches the discarded data. This is the significant bandwith waste. In connection with the fact that storage read blocks all writers it is a real problem


https://perconadev.atlassian.net/browse/PBS-42
Implemented basics for sending binlog events in position-based mode.
Current limitations: