From 4b753f1b733e2c78966e39564123805cce7516e6 Mon Sep 17 00:00:00 2001 From: Kamil Holubicki Date: Fri, 2 Oct 2026 10:10:51 +0200 Subject: [PATCH 1/3] PBS-53 feature: support blocking mode for COM_BINLOG_DUMP https://perconadev.atlassian.net/browse/PBS-53 When the client's COM_BINLOG_DUMP packet does not carry the BINLOG_DUMP_NON_BLOCK flag, the server now waits for new events to appear instead of disconnecting after serving the last available event. The current implementation uses a simple polling loop with a hardcoded 500 ms interval. Co-Authored-By: Claude Opus 4.7 --- src/minimysql/network_service.cpp | 34 ++++++++++++++++++++++++------- src/operations/sender_context.cpp | 10 ++++----- 2 files changed, 32 insertions(+), 12 deletions(-) diff --git a/src/minimysql/network_service.cpp b/src/minimysql/network_service.cpp index 30a566a..daedb6b 100644 --- a/src/minimysql/network_service.cpp +++ b/src/minimysql/network_service.cpp @@ -52,6 +52,8 @@ #pragma GCC diagnostic pop #include +#include +#include #include #include @@ -518,14 +520,36 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { } break; case minimysql::client_command_type::binlog_dump: { static constexpr auto block_size{1048576UZ}; + static constexpr auto blocking_poll_interval{ + std::chrono::milliseconds{500}}; + + const bool blocking{!context.check_binlog_non_blocking_dump()}; // TODO: initialize sender_context with binlog_name:position extracted // from the COM_BINLOG_DUMP command. operations::sender_context sender_ctx{logger, storage, block_size}; + boost::asio::steady_timer idle_timer{ + co_await boost::asio::this_coro::executor}; bool fetch_result{}; util::const_byte_span event_data{}; - while ((fetch_result = sender_ctx.get_event(event_data)) && - !event_data.empty()) { + for (;;) { + fetch_result = sender_ctx.get_event(event_data); + if (!fetch_result) { + logger->log_format( + binsrv::log_severity::error, + "net : failed to fetch next event block for {}", + remote_endpoint_str); + terminated = true; + break; + } + if (event_data.empty()) { + if (!blocking) { + break; + } + idle_timer.expires_after(blocking_poll_interval); + co_await idle_timer.async_wait(boost::asio::use_awaitable); + continue; + } const auto event_frame{ context.generate_encoded_binlog_event(event_data)}; print_generic(*logger, remote_endpoint, context, "binlog event"); @@ -536,11 +560,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { "net : sent server binlog event ({} bytes to {})", 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; + if (terminated) { break; } diff --git a/src/operations/sender_context.cpp b/src/operations/sender_context.cpp index 0050765..c1a4ef1 100644 --- a/src/operations/sender_context.cpp +++ b/src/operations/sender_context.cpp @@ -60,11 +60,11 @@ sender_context::~sender_context() = default; } 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; + // On EOF, 'storage::fetch_event_block()' leaves 'range_' with length 0, + // and unchanged offset. Restore the length while + // keeping 'binlog_name_' and the current offset so that a subsequent + // call resumes polling at the same position. + range_ = util::byte_range{range_.get_offset(), block_size_}; // setting the event span to an empty object to indicate EOF event = util::const_byte_span{}; From 35bf1429026af08ec5efafd8690dd2c79155b89c Mon Sep 17 00:00:00 2001 From: Kamil Holubicki Date: Fri, 2 Oct 2026 12:04:51 +0200 Subject: [PATCH 2/3] PBS-53 test: add MTR coverage for blocking COM_BINLOG_DUMP https://perconadev.atlassian.net/browse/PBS-53 Introduced 'binlog_streaming.rs_blocking_dump' which verifies that PBS keeps a COM_BINLOG_DUMP connection open and delivers newly produced events to downstream clients when the request does not carry the BINLOG_DUMP_NON_BLOCK flag. Flow: * Start 'binlog_server pull' in background with a 1-second checkpoint interval so small test transactions flush to storage quickly. * Spawn two concurrent 'mysqlbinlog --read-from-remote-server --stop-never' sessions against PBS. '--stop-never' clears BINLOG_DUMP_NON_BLOCK, so PBS must enter the blocking path. * Execute a first batch of transactions on the source and wait for the resulting DDL to appear in both downstream output files. * Sleep for 3 seconds with the source quiet to exercise the per-session idle-polling loop. * Execute a second batch and wait for its DDL to appear in both downstream output files. Without blocking-mode support, the sessions would have been closed on EOF after the first batch and the second wait would time out, failing the test. Co-Authored-By: Claude Opus 4.7 --- .../r/rs_blocking_dump.result | 80 +++++++++ mtr/binlog_streaming/t/rs_blocking_dump.test | 153 ++++++++++++++++++ 2 files changed, 233 insertions(+) create mode 100644 mtr/binlog_streaming/r/rs_blocking_dump.result create mode 100644 mtr/binlog_streaming/t/rs_blocking_dump.test diff --git a/mtr/binlog_streaming/r/rs_blocking_dump.result b/mtr/binlog_streaming/r/rs_blocking_dump.result new file mode 100644 index 0000000..3270e8f --- /dev/null +++ b/mtr/binlog_streaming/r/rs_blocking_dump.result @@ -0,0 +1,80 @@ +*** Resetting replication at the very beginning of the test. + +*** Generating a configuration file in JSON format for the Binlog +*** Server utility. + +*** Determining binlog file directory from the server. + +*** Creating a temporary directory for storing +*** binlog files downloaded via the Binlog Server utility. + +*** Determining the first binary log name. + +*** Starting Binlog Server Utility in background in pull mode. +include/read_file_to_var.inc + +*** Waiting for the replication source listener port ready to accept +*** connections. + +*** Starting the first downstream mysqlbinlog --read-from-remote-server +*** --stop-never session against PBS. --stop-never clears the +*** BINLOG_DUMP_NON_BLOCK flag on the COM_BINLOG_DUMP packet, so PBS +*** must keep the connection open and deliver events as they become +*** available instead of disconnecting on EOF. +include/read_file_to_var.inc + +*** Starting a second downstream mysqlbinlog session against PBS so +*** that two concurrent blocking COM_BINLOG_DUMP sessions are +*** exercised at the same time. +include/read_file_to_var.inc + +*** First batch of transactions on the source. +CREATE DATABASE test1; +CREATE TABLE test1.t(id INT UNSIGNED NOT NULL PRIMARY KEY) ENGINE=InnoDB; +INSERT INTO test1.t VALUES (1), (2), (3); + +*** Waiting for the first batch to reach the first downstream +*** mysqlbinlog through its blocking COM_BINLOG_DUMP connection. +include/wait_for_pattern_in_file.inc [test1] + +*** Waiting for the first batch to reach the second downstream +*** mysqlbinlog as well. +include/wait_for_pattern_in_file.inc [test1] + +*** Idle period: no events are produced on the source for a few +*** seconds. Both blocking connections must remain open throughout. + +*** Second batch of transactions on the source. If blocking-mode +*** support were missing, the connections would have closed on EOF +*** after the first batch and these events would never reach the +*** downstreams (the following waits would then time out). +CREATE DATABASE test2; +CREATE TABLE test2.t(id INT UNSIGNED NOT NULL PRIMARY KEY) ENGINE=InnoDB; +INSERT INTO test2.t VALUES (4), (5), (6); + +*** Waiting for the second batch to reach the first downstream +*** mysqlbinlog. +include/wait_for_pattern_in_file.inc [test2] + +*** Waiting for the second batch to reach the second downstream +*** mysqlbinlog. +include/wait_for_pattern_in_file.inc [test2] + +*** Sending SIGTERM to the first downstream mysqlbinlog session and +*** waiting for it to terminate. + +*** Sending SIGTERM to the second downstream mysqlbinlog session and +*** waiting for it to terminate. + +*** Dropping the test databases. +DROP DATABASE test1; +DROP DATABASE test2; + +*** Sending SIGTERM signal to the Binlog Server Utility and waiting +*** for the process to terminate. + +*** Removing the Binlog Server utility storage directory. + +*** Removing the Binlog Server utility log file. + +*** Removing the Binlog Server utility configuration file. diff --git a/mtr/binlog_streaming/t/rs_blocking_dump.test b/mtr/binlog_streaming/t/rs_blocking_dump.test new file mode 100644 index 0000000..0a5c4a7 --- /dev/null +++ b/mtr/binlog_streaming/t/rs_blocking_dump.test @@ -0,0 +1,153 @@ +--source ../include/have_binsrv.inc + +--source ../include/v80_v84_compatibility_defines.inc + +--source include/count_sessions.inc + +# in case of --repeat=N, we need to start from a fresh binary log to make +# this test deterministic +--echo *** Resetting replication at the very beginning of the test. +--disable_query_log +eval $stmt_reset_binary_logs_and_gtids; +--enable_query_log + +# identifying backend storage type ('file' or 's3') +--source ../include/identify_storage_backend.inc + +# A 1-second checkpoint interval makes PBS flush small test transactions +# to storage quickly, so the downstream mysqlbinlog sessions can observe +# them within a few seconds. 'binsrv_idle_time' controls how long PBS +# waits before re-connecting to the upstream source after a read timeout. +--let $binsrv_connect_timeout = 3 +--let $binsrv_read_timeout = 3 +--let $binsrv_idle_time = 1 +--let $binsrv_verify_checksum = TRUE +--let $binsrv_replication_mode = position +--let $binsrv_checkpoint_interval = 1s +--source ../include/set_up_binsrv_environment.inc + +--echo +--echo *** Determining the first binary log name. +--let $first_binlog = query_get_value($stmt_show_binary_log_status, File, 1) + +--echo +--echo *** Starting Binlog Server Utility in background in pull mode. +--let $proc_command_line = $BINSRV pull $binsrv_config_file_path > /dev/null +--source ../include/start_proc_in_background.inc +--let $binsrv_pid = $proc_pid + +--echo +--echo *** Waiting for the replication source listener port ready to accept +--echo *** connections. +--let $listening_port = $binsrv_replication_source_port +--source ../include/wait_for_listening_port.inc + +--echo +--echo *** Starting the first downstream mysqlbinlog --read-from-remote-server +--echo *** --stop-never session against PBS. --stop-never clears the +--echo *** BINLOG_DUMP_NON_BLOCK flag on the COM_BINLOG_DUMP packet, so PBS +--echo *** must keep the connection open and deliver events as they become +--echo *** available instead of disconnecting on EOF. +--let $downstream1_dump = $MYSQL_TMP_DIR/rs_blocking_dump_1.sql +--let $proc_command_line = $MYSQL_BINLOG --read-from-remote-server --stop-never --host=127.0.0.1 --port=$binsrv_replication_source_port --user=$binsrv_auth_user --password=$binsrv_auth_password --default-auth=caching_sha2_password $first_binlog > $downstream1_dump 2>/dev/null +--source ../include/start_proc_in_background.inc +--let $downstream1_pid = $proc_pid + +--echo +--echo *** Starting a second downstream mysqlbinlog session against PBS so +--echo *** that two concurrent blocking COM_BINLOG_DUMP sessions are +--echo *** exercised at the same time. +--let $downstream2_dump = $MYSQL_TMP_DIR/rs_blocking_dump_2.sql +--let $proc_command_line = $MYSQL_BINLOG --read-from-remote-server --stop-never --host=127.0.0.1 --port=$binsrv_replication_source_port --user=$binsrv_auth_user --password=$binsrv_auth_password --default-auth=caching_sha2_password $first_binlog > $downstream2_dump 2>/dev/null +--source ../include/start_proc_in_background.inc +--let $downstream2_pid = $proc_pid + +--echo +--echo *** First batch of transactions on the source. +CREATE DATABASE test1; +CREATE TABLE test1.t(id INT UNSIGNED NOT NULL PRIMARY KEY) ENGINE=InnoDB; +INSERT INTO test1.t VALUES (1), (2), (3); + +--echo +--echo *** Waiting for the first batch to reach the first downstream +--echo *** mysqlbinlog through its blocking COM_BINLOG_DUMP connection. +--let $grep_pattern = test1 +--let $grep_file = $downstream1_dump +--let $wait_timeout = 60 +--source include/wait_for_pattern_in_file.inc + +--echo +--echo *** Waiting for the first batch to reach the second downstream +--echo *** mysqlbinlog as well. +--let $grep_pattern = test1 +--let $grep_file = $downstream2_dump +--let $wait_timeout = 60 +--source include/wait_for_pattern_in_file.inc + +--echo +--echo *** Idle period: no events are produced on the source for a few +--echo *** seconds. Both blocking connections must remain open throughout. +--sleep 3 + +--echo +--echo *** Second batch of transactions on the source. If blocking-mode +--echo *** support were missing, the connections would have closed on EOF +--echo *** after the first batch and these events would never reach the +--echo *** downstreams (the following waits would then time out). +CREATE DATABASE test2; +CREATE TABLE test2.t(id INT UNSIGNED NOT NULL PRIMARY KEY) ENGINE=InnoDB; +INSERT INTO test2.t VALUES (4), (5), (6); + +--echo +--echo *** Waiting for the second batch to reach the first downstream +--echo *** mysqlbinlog. +--let $grep_pattern = test2 +--let $grep_file = $downstream1_dump +--let $wait_timeout = 60 +--source include/wait_for_pattern_in_file.inc + +--echo +--echo *** Waiting for the second batch to reach the second downstream +--echo *** mysqlbinlog. +--let $grep_pattern = test2 +--let $grep_file = $downstream2_dump +--let $wait_timeout = 60 +--source include/wait_for_pattern_in_file.inc + +--echo +--echo *** Sending SIGTERM to the first downstream mysqlbinlog session and +--echo *** waiting for it to terminate. +--let $proc_pid = $downstream1_pid +--source ../include/terminate_proc.inc + +--echo +--echo *** Sending SIGTERM to the second downstream mysqlbinlog session and +--echo *** waiting for it to terminate. +--let $proc_pid = $downstream2_pid +--source ../include/terminate_proc.inc + +--echo +--echo *** Dropping the test databases. +DROP DATABASE test1; +DROP DATABASE test2; + +--echo +--echo *** Sending SIGTERM signal to the Binlog Server Utility and waiting +--echo *** for the process to terminate. +--let $proc_pid = $binsrv_pid +--source ../include/terminate_proc.inc + +--remove_file $downstream1_dump +--remove_file $downstream2_dump + +# cleaning up +--source ../include/tear_down_binsrv_environment.inc + +# As the Binlog Server Utility interrupts the connection upon timeout, here we +# need to close it on the MySQL server side as well in order to make sure that +# MTR 'check-test' before and after the test produces the same output. +--source ../include/kill_binlog_dump_connection.inc + +# Also, we use 'count_sessions' include files to make sure that 'Binlog Dump' +# connection is indeed closed. +--source include/wait_until_count_sessions.inc From 675b1de47c4fb9950170b4c5230f2e5b4cce250f Mon Sep 17 00:00:00 2001 From: Kamil Holubicki Date: Fri, 2 Oct 2026 14:03:37 +0200 Subject: [PATCH 3/3] PBS-53 fixup: extract COM_BINLOG_DUMP handler to tame session() complexity https://perconadev.atlassian.net/browse/PBS-53 The blocking-mode branches added to the 'binlog_dump' case in 'minimysql::session()' pushed the function over clang-tidy's 'readability-function-cognitive-complexity' threshold (30 vs 25) and broke the Clang 20 CI jobs. Moved the body of the 'binlog_dump' case into a new free coroutine 'handle_binlog_dump_command()' in the anonymous namespace. The case body now reduces to one 'co_await' plus setting 'terminated = true', which brings 'session()' back under the limit. Behavior is unchanged. Co-Authored-By: Claude Opus 4.7 --- src/minimysql/network_service.cpp | 111 ++++++++++++++++-------------- 1 file changed, 60 insertions(+), 51 deletions(-) diff --git a/src/minimysql/network_service.cpp b/src/minimysql/network_service.cpp index daedb6b..7d1a6a4 100644 --- a/src/minimysql/network_service.cpp +++ b/src/minimysql/network_service.cpp @@ -276,6 +276,63 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { } } +[[nodiscard]] boost::asio::awaitable handle_binlog_dump_command( + // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) + const binsrv::basic_logger_ptr &logger, + // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) + const binsrv::storage_ptr &storage, + // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) + boost::asio::ip::tcp::socket &socket, + // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) + minimysql::connection_context &context, + // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) + const boost::asio::ip::tcp::endpoint &remote_endpoint, + // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) + const std::string &remote_endpoint_str, + std::chrono::seconds write_timeout) { + static constexpr auto block_size{1048576UZ}; + static constexpr auto blocking_poll_interval{std::chrono::milliseconds{500}}; + + const bool blocking{!context.check_binlog_non_blocking_dump()}; + + // TODO: initialize sender_context with binlog_name:position extracted + // from the COM_BINLOG_DUMP command. + operations::sender_context sender_ctx{logger, storage, block_size}; + boost::asio::steady_timer idle_timer{ + co_await boost::asio::this_coro::executor}; + util::const_byte_span event_data{}; + for (;;) { + if (!sender_ctx.get_event(event_data)) { + logger->log_format(binsrv::log_severity::error, + "net : failed to fetch next event block for {}", + remote_endpoint_str); + co_return; + } + if (event_data.empty()) { + if (!blocking) { + break; + } + idle_timer.expires_after(blocking_poll_interval); + co_await idle_timer.async_wait(boost::asio::use_awaitable); + continue; + } + 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(binsrv::log_severity::debug, + "net : sent server binlog event ({} bytes to {})", + std::size(event_frame), remote_endpoint_str); + } + + const auto eof = context.generate_encoded_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); +} + #pragma GCC diagnostic push #pragma GCC diagnostic ignored "-Wmismatched-new-delete" @@ -519,57 +576,9 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { std::size(ok_after_ping), remote_endpoint_str); } break; case minimysql::client_command_type::binlog_dump: { - static constexpr auto block_size{1048576UZ}; - static constexpr auto blocking_poll_interval{ - std::chrono::milliseconds{500}}; - - const bool blocking{!context.check_binlog_non_blocking_dump()}; - - // TODO: initialize sender_context with binlog_name:position extracted - // from the COM_BINLOG_DUMP command. - operations::sender_context sender_ctx{logger, storage, block_size}; - boost::asio::steady_timer idle_timer{ - co_await boost::asio::this_coro::executor}; - bool fetch_result{}; - util::const_byte_span event_data{}; - for (;;) { - fetch_result = sender_ctx.get_event(event_data); - if (!fetch_result) { - logger->log_format( - binsrv::log_severity::error, - "net : failed to fetch next event block for {}", - remote_endpoint_str); - terminated = true; - break; - } - if (event_data.empty()) { - if (!blocking) { - break; - } - idle_timer.expires_after(blocking_poll_interval); - co_await idle_timer.async_wait(boost::asio::use_awaitable); - continue; - } - 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( - binsrv::log_severity::debug, - "net : sent server binlog event ({} bytes to {})", - std::size(event_frame), remote_endpoint_str); - } - if (terminated) { - break; - } - - const auto eof = context.generate_encoded_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); + co_await handle_binlog_dump_command(logger, storage, socket, context, + remote_endpoint, + remote_endpoint_str, write_timeout); terminated = true; } break; case minimysql::client_command_type::quit: {