From 31fb5c6a1548ba7df88e93eebe0c79f55fa0f36e Mon Sep 17 00:00:00 2001 From: Kamil Holubicki Date: Fri, 18 Sep 2026 17:01:33 +0200 Subject: [PATCH 1/3] PBS-3 feature: add replication_source config section for pull mode listener Introduce a new top-level "replication_source" JSON configuration section that controls the MySQL-compatible listener the utility exposes in 'pull' mode. The section carries three fields: - port : TCP port the listener binds to - read_timeout : per-frame read timeout applied to accepted sessions - write_timeout : per-frame write timeout applied to accepted sessions replication_source is a required section validated at config load time (port, read_timeout and write_timeout must all be non-zero). The README and the sample main_config.json are updated accordingly. Co-Authored-By: Claude Opus 4.7 --- CMakeLists.txt | 4 + README.md | 11 +++ main_config.json | 5 ++ src/binsrv/main_config.cpp | 1 + src/binsrv/main_config.hpp | 20 ++--- src/binsrv/replication_source_config.cpp | 42 ++++++++++ src/binsrv/replication_source_config.hpp | 40 ++++++++++ src/binsrv/replication_source_config_fwd.hpp | 25 ++++++ src/minimysql/network_service.cpp | 82 ++++++++++---------- src/minimysql/network_service.hpp | 8 +- src/operations/pull_operation.cpp | 12 +-- 11 files changed, 190 insertions(+), 60 deletions(-) create mode 100644 src/binsrv/replication_source_config.cpp create mode 100644 src/binsrv/replication_source_config.hpp create mode 100644 src/binsrv/replication_source_config_fwd.hpp diff --git a/CMakeLists.txt b/CMakeLists.txt index 65f2555..da7c947 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -600,6 +600,10 @@ set(binsrv_source_files src/binsrv/replication_mode_type_fwd.hpp src/binsrv/replication_mode_type.hpp + src/binsrv/replication_source_config_fwd.hpp + src/binsrv/replication_source_config.hpp + src/binsrv/replication_source_config.cpp + src/binsrv/rewrite_config_fwd.hpp src/binsrv/rewrite_config.hpp src/binsrv/rewrite_config.cpp diff --git a/README.md b/README.md index 588671d..c4960ae 100644 --- a/README.md +++ b/README.md @@ -514,6 +514,11 @@ The Percona Binary Log Server configuration file has the following format. "file_size": "128M" } }, + "replication_source": { + "port": 3307, + "read_timeout": 60, + "write_timeout": 60 + }, "keyring": { "uri": "file:///var/lib/pbs/keyring/keyring_data.json" }, @@ -584,6 +589,12 @@ If this section is present, then the utility will not split binlog events the sa - `` - the base name of the generated binlog file names in the "rewrite" mode. E.g. `rewritten_binlog` will cause `rewritten_binlog.000001`, `rewritten_binlog.000002`, etc. file names to be generated. - `` - the maximum individual binlog file size after reaching which the utility will switch to a new one. The value is expected to be a string containing an integer followed by an optional suffix 'K' / 'M' / 'G' / 'T' / 'P', e.g. /\d+\[KMGTP\]?/. The minimal allowed value of this parameter is `1024` bytes. +#### \ section +This section configures the built-in MySQL-compatible listener the utility exposes in `pull` mode so that downstream replicas can dump binary log events from it (the utility acts as a replication source). +- `` - the TCP port on which the utility listens for incoming replica connections. +- `` - the number of seconds the utility will wait to read data from a connected replica before treating the connection as timed out. +- `` - the number of seconds the utility will wait to write data to a connected replica before treating the connection as timed out. + #### \ section If this an optional section that specifies keyring configuration parameters. It must be present if the storage has at least one encrypted binlog file. - `` - specifies location of the keyring JSON data file (currently only 'file://' scheme is supported meaning that the file should be taken from the local file sytem from the path specified in this URI, e.g. `file:///var/lib/pbs/keyring/keyring_data.json`). diff --git a/main_config.json b/main_config.json index 28da1b5..db69b9e 100644 --- a/main_config.json +++ b/main_config.json @@ -36,6 +36,11 @@ "file_size": "128M" } }, + "replication_source": { + "port": 3307, + "read_timeout": 60, + "write_timeout": 60 + }, "keyring": { "uri": "file:///home/user/keyring/keyring/keyring_data.json" }, diff --git a/src/binsrv/main_config.cpp b/src/binsrv/main_config.cpp index 95e5682..9406e2e 100644 --- a/src/binsrv/main_config.cpp +++ b/src/binsrv/main_config.cpp @@ -60,6 +60,7 @@ void main_config::validate() const { root().get<"connection">().validate(); root().get<"storage">().validate(); root().get<"replication">().validate(); + root().get<"replication_source">().validate(); } } // namespace binsrv diff --git a/src/binsrv/main_config.hpp b/src/binsrv/main_config.hpp index bdd0572..f1d2288 100644 --- a/src/binsrv/main_config.hpp +++ b/src/binsrv/main_config.hpp @@ -18,10 +18,11 @@ #include "binsrv/main_config_fwd.hpp" // IWYU pragma: export -#include "binsrv/keyring_config.hpp" // IWYU pragma: export -#include "binsrv/logger_config.hpp" // IWYU pragma: export -#include "binsrv/replication_config.hpp" // IWYU pragma: export -#include "binsrv/storage_config.hpp" // IWYU pragma: export +#include "binsrv/keyring_config.hpp" // IWYU pragma: export +#include "binsrv/logger_config.hpp" // IWYU pragma: export +#include "binsrv/replication_config.hpp" // IWYU pragma: export +#include "binsrv/replication_source_config.hpp" // IWYU pragma: export +#include "binsrv/storage_config.hpp" // IWYU pragma: export #include "easymysql/connection_config.hpp" // IWYU pragma: export @@ -33,11 +34,12 @@ class [[nodiscard]] main_config { private: using impl_type = util::nv_tuple< // clang-format off - util::nv<"logger" , logger_config>, - util::nv<"connection" , easymysql::connection_config>, - util::nv<"replication", binsrv::replication_config>, - util::nv<"keyring" , optional_keyring_config>, - util::nv<"storage" , storage_config> + util::nv<"logger" , logger_config>, + util::nv<"connection" , easymysql::connection_config>, + util::nv<"replication" , binsrv::replication_config>, + util::nv<"replication_source", binsrv::replication_source_config>, + util::nv<"keyring" , optional_keyring_config>, + util::nv<"storage" , storage_config> // clang-format on >; diff --git a/src/binsrv/replication_source_config.cpp b/src/binsrv/replication_source_config.cpp new file mode 100644 index 0000000..2e4b5a2 --- /dev/null +++ b/src/binsrv/replication_source_config.cpp @@ -0,0 +1,42 @@ +// 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 "binsrv/replication_source_config.hpp" + +#include + +#include "util/exception_location_helpers.hpp" + +namespace binsrv { + +void replication_source_config::validate() const { + if (get<"port">() == 0U) { + util::exception_location().raise( + "error validating replication source config: " + "port must be greater than 0"); + } + if (get<"read_timeout">() == 0U) { + util::exception_location().raise( + "error validating replication source config: " + "read_timeout must be greater than 0"); + } + if (get<"write_timeout">() == 0U) { + util::exception_location().raise( + "error validating replication source config: " + "write_timeout must be greater than 0"); + } +} + +} // namespace binsrv diff --git a/src/binsrv/replication_source_config.hpp b/src/binsrv/replication_source_config.hpp new file mode 100644 index 0000000..b1dba57 --- /dev/null +++ b/src/binsrv/replication_source_config.hpp @@ -0,0 +1,40 @@ +// 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 BINSRV_REPLICATION_SOURCE_CONFIG_HPP +#define BINSRV_REPLICATION_SOURCE_CONFIG_HPP + +#include "binsrv/replication_source_config_fwd.hpp" // IWYU pragma: export + +#include + +#include "util/nv_tuple.hpp" + +namespace binsrv { + +struct [[nodiscard]] replication_source_config + : util::nv_tuple< + // clang-format off + util::nv<"port" , std::uint16_t>, + util::nv<"read_timeout" , std::uint32_t>, + util::nv<"write_timeout", std::uint32_t> + // clang-format on + > { + void validate() const; +}; + +} // namespace binsrv + +#endif // BINSRV_REPLICATION_SOURCE_CONFIG_HPP diff --git a/src/binsrv/replication_source_config_fwd.hpp b/src/binsrv/replication_source_config_fwd.hpp new file mode 100644 index 0000000..abc4117 --- /dev/null +++ b/src/binsrv/replication_source_config_fwd.hpp @@ -0,0 +1,25 @@ +// 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 BINSRV_REPLICATION_SOURCE_CONFIG_FWD_HPP +#define BINSRV_REPLICATION_SOURCE_CONFIG_FWD_HPP + +namespace binsrv { + +struct replication_source_config; + +} // namespace binsrv + +#endif // BINSRV_REPLICATION_SOURCE_CONFIG_FWD_HPP diff --git a/src/minimysql/network_service.cpp b/src/minimysql/network_service.cpp index 7e9acf6..1658a1e 100644 --- a/src/minimysql/network_service.cpp +++ b/src/minimysql/network_service.cpp @@ -281,6 +281,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { binsrv::basic_logger &logger, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) binsrv::storage &storage, boost::asio::ip::tcp::socket socket, + // NOLINTNEXTLINE(bugprone-easily-swappable-parameters) + std::chrono::seconds read_timeout, std::chrono::seconds write_timeout, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) const std::string &username, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) @@ -310,9 +312,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { const auto server_greeting{context.generate_encoded_server_greeting()}; print_server_greeting(logger, remote_endpoint, context); - co_await minimysql::async_write_mysql_frame( - socket, server_greeting, - network_service::session_authentication_timeout); + 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); @@ -327,8 +328,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // auth_method_name // attributes - co_await minimysql::async_read_mysql_frame( - socket, data, network_service::session_authentication_timeout); + 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); @@ -340,9 +340,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { "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"); - co_await minimysql::async_write_mysql_frame( - socket, access_denied, - network_service::session_authentication_timeout); + 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); @@ -360,16 +359,14 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { const auto auth_method_switch{ context.generate_encoded_auth_method_switch()}; print_generic(logger, remote_endpoint, context, "auth method switch"); - co_await minimysql::async_write_mysql_frame( - socket, auth_method_switch, - network_service::session_authentication_timeout); + co_await minimysql::async_write_mysql_frame(socket, auth_method_switch, + write_timeout); 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, network_service::session_authentication_timeout); + 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 {})", @@ -383,9 +380,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { context.get_client_username()); const auto access_denied{context.generate_encoded_access_denied()}; print_error(logger, remote_endpoint, context, "auth failure"); - co_await minimysql::async_write_mysql_frame( - socket, access_denied, - network_service::session_authentication_timeout); + 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); @@ -400,9 +396,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { const auto fast_auth_success{context.generate_encoded_fast_auth()}; print_generic(logger, remote_endpoint, context, "auth method data (fast auth)"); - co_await minimysql::async_write_mysql_frame( - socket, fast_auth_success, - network_service::session_authentication_timeout); + 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); @@ -410,8 +405,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // sending server ok after successful authentication const auto auth_ok{context.generate_encoded_ok()}; print_generic(logger, remote_endpoint, context, "ok (auth)"); - co_await minimysql::async_write_mysql_frame( - socket, auth_ok, network_service::session_authentication_timeout); + co_await minimysql::async_write_mysql_frame(socket, auth_ok, write_timeout); logger.log_format( binsrv::log_severity::debug, "net : sent server ok after authentication ({} bytes to {})", @@ -476,8 +470,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { bool terminated{false}; while (!terminated) { context.enter_command_loop_iteration(); - co_await minimysql::async_read_mysql_frame( - socket, data, network_service::session_command_timeout); + 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); @@ -491,8 +484,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { if (known_query_it != std::end(known_queries)) { const auto resultset{known_query_it->second(context)}; print_generic(logger, remote_endpoint, context, "resultset"); - co_await minimysql::async_write_mysql_frames( - socket, resultset, network_service::session_command_timeout); + 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); @@ -500,8 +493,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { // return 'syntax error' for every other query const auto syntax_error = context.generate_encoded_syntax_error(); print_error(logger, remote_endpoint, context, "syntax error"); - co_await minimysql::async_write_mysql_frame( - socket, syntax_error, network_service::session_command_timeout); + co_await minimysql::async_write_mysql_frame(socket, syntax_error, + write_timeout); logger.log_format( binsrv::log_severity::debug, "net : sent server syntax error ({} bytes to {})", @@ -511,8 +504,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { case minimysql::client_command_type::ping: { const auto ok_after_ping{context.generate_encoded_ok()}; print_generic(logger, remote_endpoint, context, "ok (ping success)"); - co_await minimysql::async_write_mysql_frame( - socket, ok_after_ping, network_service::session_command_timeout); + 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); @@ -525,8 +518,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { 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, network_service::session_command_timeout); + co_await minimysql::async_write_mysql_frame(socket, event, + write_timeout); logger.log_format( binsrv::log_severity::debug, "net : sent server binlog event ({} bytes to {})", @@ -534,8 +527,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { } const auto eof = context.generate_encoded_eof(); print_generic(logger, remote_endpoint, context, "binlog eof"); - co_await minimysql::async_write_mysql_frame( - socket, eof, network_service::session_command_timeout); + 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); @@ -551,8 +543,7 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { context.generate_encoded_unknown_command(); print_error(logger, remote_endpoint, context, "unknown command"); co_await minimysql::async_write_mysql_frame( - socket, unknown_command_error, - network_service::session_command_timeout); + socket, unknown_command_error, write_timeout); logger.log_format( binsrv::log_severity::debug, "net : sent server unknown command ({} bytes to {})", @@ -575,6 +566,8 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { binsrv::storage &storage, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) boost::asio::ip::tcp::acceptor &acceptor, + // NOLINTNEXTLINE(bugprone-easily-swappable-parameters) + std::chrono::seconds read_timeout, std::chrono::seconds write_timeout, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) const std::string &username, // NOLINTNEXTLINE(cppcoreguidelines-avoid-reference-coroutine-parameters) @@ -603,10 +596,11 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) { boost::lexical_cast(remote_endpoint)); // NOLINTNEXTLINE(misc-include-cleaner) - boost::asio::co_spawn( - executor, - session(logger, storage, std::move(socket), username, password), - boost::asio::detached); + boost::asio::co_spawn(executor, + session(logger, storage, std::move(socket), + read_timeout, write_timeout, username, + password), + boost::asio::detached); } } catch (...) { handle_exception(logger, "listener"); @@ -619,6 +613,8 @@ network_service::network_service( binsrv::basic_logger_ptr logger, boost::asio::io_context &context, binsrv::storage_ptr storage, std::uint16_t listening_port, // NOLINTNEXTLINE(bugprone-easily-swappable-parameters) + std::chrono::seconds read_timeout, std::chrono::seconds write_timeout, + // NOLINTNEXTLINE(bugprone-easily-swappable-parameters) std::string_view username, std::string_view password) : logger_{std::move(logger)}, storage_{std::move(storage)}, username_(username), password_(password), context_{&context}, @@ -627,10 +623,10 @@ network_service::network_service( listening_port})} { assert(logger_); // NOLINTNEXTLINE(misc-include-cleaner) - boost::asio::co_spawn( - *context_, - listener(*logger_, *storage_, *acceptor_, username_, password_), - boost::asio::detached); + boost::asio::co_spawn(*context_, + listener(*logger_, *storage_, *acceptor_, read_timeout, + write_timeout, username_, password_), + boost::asio::detached); } network_service::~network_service() = default; diff --git a/src/minimysql/network_service.hpp b/src/minimysql/network_service.hpp index 85ca5e5..b4dae22 100644 --- a/src/minimysql/network_service.hpp +++ b/src/minimysql/network_service.hpp @@ -16,7 +16,9 @@ #ifndef MINIMYSQL_NETWORK_SERVICE_HPP #define MINIMYSQL_NETWORK_SERVICE_HPP +#include #include +#include #include #include @@ -29,12 +31,12 @@ namespace minimysql { class network_service { public: static constexpr auto expected_packet_size{4096UZ}; - static constexpr std::chrono::seconds session_authentication_timeout{10}; - static constexpr std::chrono::seconds session_command_timeout{120}; network_service(binsrv::basic_logger_ptr logger, boost::asio::io_context &context, binsrv::storage_ptr storage, - std::uint16_t listening_port, std::string_view username, + std::uint16_t listening_port, + std::chrono::seconds read_timeout, + std::chrono::seconds write_timeout, std::string_view username, std::string_view password); network_service(const network_service &) = delete; diff --git a/src/operations/pull_operation.cpp b/src/operations/pull_operation.cpp index c28b66f..e51d120 100644 --- a/src/operations/pull_operation.cpp +++ b/src/operations/pull_operation.cpp @@ -82,8 +82,6 @@ generic_operation::generic_operation( : basic_operation{cmd_args, expected_number_of_arguments} {} [[nodiscard]] bool generic_operation::execute() const { - static constexpr std::uint16_t listening_port{3307}; - static constexpr std::string_view default_username{"rpl"}; static constexpr std::string_view default_password{"password"}; @@ -119,9 +117,13 @@ generic_operation::generic_operation( easymysql::connection_replication_mode_type::blocking, config, logger, storage}; - const minimysql::network_service service(logger, io_ctx, storage, - listening_port, default_username, - default_password); + const auto &replication_source_config{ + config->root().get<"replication_source">()}; + const minimysql::network_service service( + logger, io_ctx, storage, replication_source_config.get<"port">(), + std::chrono::seconds{replication_source_config.get<"read_timeout">()}, + std::chrono::seconds{replication_source_config.get<"write_timeout">()}, + default_username, default_password); const auto idle_time_seconds{ config->root().get<"replication">().get<"idle_time">()}; From 98c51594f240f8995cd4dad58c4945003ffdf990 Mon Sep 17 00:00:00 2001 From: Kamil Holubicki Date: Fri, 18 Sep 2026 17:03:55 +0200 Subject: [PATCH 2/3] PBS-3 test: run pull-mode MTR tests in parallel with per-worker ports Now that the Binlog Server pull-mode listener reads its TCP port from the replication_source config section instead of hard-coding 3307, generate_binsrv_config.inc emits a replication_source block whose port is @@global.port + 1000. Because @@global.port is unique per MTR worker, so is the derived listener port, so pull-mode tests no longer collide when run in parallel. The +1000 offset intentionally jumps well beyond MTR's per-worker port allocation (empirically MTR reserves ports within a few hundred of each worker's @@port on GitHub Actions runners), so the derived listener port never lands inside a sibling worker's reserved range. The computed value is exported to callers via $binsrv_replication_source_port for tests that need to talk to the listener themselves. Consequently: - pull_mode.test, binlog_flush.test and auth_method_switch.test drop their '--source include/not_parallel.inc' guards. - auth_method_switch.test's /dev/tcp probe and the two mysql CLI invocations that connect to the listener now use $binsrv_replication_source_port instead of hard-coded 3307; the matching .result echo is rewritten in port-agnostic form. Co-Authored-By: Claude Opus 4.7 --- .../include/generate_binsrv_config.inc | 24 +++++++++++++++++++ .../r/auth_method_switch.result | 10 ++++---- .../t/auth_method_switch.test | 22 +++++++---------- mtr/binlog_streaming/t/binlog_flush.test | 4 ---- mtr/binlog_streaming/t/pull_mode.test | 4 ---- 5 files changed, 38 insertions(+), 26 deletions(-) diff --git a/mtr/binlog_streaming/include/generate_binsrv_config.inc b/mtr/binlog_streaming/include/generate_binsrv_config.inc index 1c5dbe7..f1786f3 100644 --- a/mtr/binlog_streaming/include/generate_binsrv_config.inc +++ b/mtr/binlog_streaming/include/generate_binsrv_config.inc @@ -1,6 +1,11 @@ # # Creates a JSON configuration file for the Binlog Server Utility # +# Sets (as an output for the caller): +# $binsrv_replication_source_port - TCP port on which the utility +# listens for downstream replicas (@@global.port + 1000, unique per +# MTR worker). +# # Usage: # --let $binsrv_log_level = trace | debug | info | warning | error | fatal (optional, default: trace) # --let $binsrv_connection_user = repl_user (optional) @@ -88,6 +93,20 @@ if ($binsrv_connection_host == "") } } +# The Binlog Server exposes its own MySQL-compatible listener in 'pull' +# mode. To let multiple MTR workers run pull-mode tests in parallel, each +# worker gets its own listening port derived from the worker's MySQL port +# (worker port + 1000). The +1000 offset intentionally jumps well beyond +# MTR's per-worker port allocation (empirically MTR reserves ports within +# ~200 of each worker's @@port on GitHub Actions runners), so the derived +# listener port never collides with a neighbour worker's mysqld/mysqlx or +# other MTR-reserved ports. An earlier +100 offset produced bind failures +# ("Address already in use") on CI when the derived port fell inside a +# sibling worker's range. We expose the computed value via +# $binsrv_replication_source_port so the test can reuse it (e.g. when +# connecting a probe client to the listener). +--let $binsrv_replication_source_port = `SELECT @@global.port + 1000` + eval SET @binsrv_config_json = JSON_OBJECT( 'logger', JSON_OBJECT( 'level', '$binsrv_log_level', @@ -108,6 +127,11 @@ eval SET @binsrv_config_json = JSON_OBJECT( 'verify_checksum', $binsrv_verify_checksum, 'mode', '$binsrv_replication_mode' ), + 'replication_source', JSON_OBJECT( + 'port', $binsrv_replication_source_port, + 'read_timeout', 60, + 'write_timeout', 60 + ), 'storage', JSON_OBJECT( 'backend', '$storage_backend', 'uri', @storage_uri diff --git a/mtr/binlog_streaming/r/auth_method_switch.result b/mtr/binlog_streaming/r/auth_method_switch.result index 2570595..c75319a 100644 --- a/mtr/binlog_streaming/r/auth_method_switch.result +++ b/mtr/binlog_streaming/r/auth_method_switch.result @@ -17,11 +17,11 @@ include/read_file_to_var.inc *** Waiting for the Binlog Server listener to come up on -*** 127.0.0.1:3307. We probe with bash's /dev/tcp instead of the -*** mysql client because bash is not ASAN-instrumented and -*** /dev/tcp uses a plain connect(2), so each attempt is cheap and -*** measures exactly the "listening on the port" state we care -*** about. +*** 127.0.0.1:. We probe with bash's /dev/tcp +*** instead of the mysql client because bash is not +*** ASAN-instrumented and /dev/tcp uses a plain connect(2), so each +*** attempt is cheap and measures exactly the "listening on the +*** port" state we care about. *** Control: client picks the same plugin the server advertises *** (caching_sha2_password). No AuthMethodSwitch is expected on the diff --git a/mtr/binlog_streaming/t/auth_method_switch.test b/mtr/binlog_streaming/t/auth_method_switch.test index 3e9de23..1f2e378 100644 --- a/mtr/binlog_streaming/t/auth_method_switch.test +++ b/mtr/binlog_streaming/t/auth_method_switch.test @@ -1,7 +1,3 @@ -# The Binlog Server listens on a hard-coded TCP port, so this test can not -# run in parallel with other tests using the same port. ---source include/not_parallel.inc - --source ../include/have_binsrv.inc --source ../include/v80_v84_compatibility_defines.inc @@ -61,11 +57,11 @@ EOF --echo --echo *** Waiting for the Binlog Server listener to come up on ---echo *** 127.0.0.1:3307. We probe with bash's /dev/tcp instead of the ---echo *** mysql client because bash is not ASAN-instrumented and ---echo *** /dev/tcp uses a plain connect(2), so each attempt is cheap and ---echo *** measures exactly the "listening on the port" state we care ---echo *** about. +--echo *** 127.0.0.1:. We probe with bash's /dev/tcp +--echo *** instead of the mysql client because bash is not +--echo *** ASAN-instrumented and /dev/tcp uses a plain connect(2), so each +--echo *** attempt is cheap and measures exactly the "listening on the +--echo *** port" state we care about. --let $max_wait = 300 --let $iteration = 0 --let $port_open = 0 @@ -74,7 +70,7 @@ while ($iteration < $max_wait) if (!$port_open) { --error 0, 1 - --exec bash -c "echo > /dev/tcp/127.0.0.1/3307" 2>/dev/null + --exec bash -c "echo > /dev/tcp/127.0.0.1/$binsrv_replication_source_port" 2>/dev/null --let $port_status = $__error if ($port_status == 0) { @@ -90,7 +86,7 @@ while ($iteration < $max_wait) } if (!$port_open) { - --die The Binlog Server listener did not become reachable on 3307 within 300 seconds + --die The Binlog Server listener did not become reachable on its port within 300 seconds } --echo @@ -98,7 +94,7 @@ if (!$port_open) --echo *** (caching_sha2_password). No AuthMethodSwitch is expected on the --echo *** wire; a zero exit code from mysql confirms the session got as --echo *** far as running the probe query. ---exec $MYSQL --protocol=TCP --host=127.0.0.1 --port=3307 --user=rpl --password=password --default-auth=caching_sha2_password --skip-column-names -e "$probe_query" >/dev/null 2>&1 +--exec $MYSQL --protocol=TCP --host=127.0.0.1 --port=$binsrv_replication_source_port --user=rpl --password=password --default-auth=caching_sha2_password --skip-column-names -e "$probe_query" >/dev/null 2>&1 --echo --echo *** Trigger: client forces mysql_native_password in its handshake so @@ -109,7 +105,7 @@ if (!$port_open) --echo *** MTR fails the --exec. The log grep that follows the shutdown --echo *** is what actually proves the switch happened - this line only --echo *** proves the session survived it. ---exec $MYSQL --protocol=TCP --host=127.0.0.1 --port=3307 --user=rpl --password=password --default-auth=mysql_native_password --skip-column-names -e "$probe_query" >/dev/null 2>&1 +--exec $MYSQL --protocol=TCP --host=127.0.0.1 --port=$binsrv_replication_source_port --user=rpl --password=password --default-auth=mysql_native_password --skip-column-names -e "$probe_query" >/dev/null 2>&1 --echo --echo *** Sending SIGTERM to the Binlog Server Utility and waiting for the diff --git a/mtr/binlog_streaming/t/binlog_flush.test b/mtr/binlog_streaming/t/binlog_flush.test index 2d297ae..005b8bc 100644 --- a/mtr/binlog_streaming/t/binlog_flush.test +++ b/mtr/binlog_streaming/t/binlog_flush.test @@ -1,7 +1,3 @@ -# temporarily marking this test as non-parallel because of the hardcoded -# listening port ---source include/not_parallel.inc - # The purpose of the test is to validate that PBS doesn't flush its internal binlog buffer # to the file every time the whole transaction is replicated. diff --git a/mtr/binlog_streaming/t/pull_mode.test b/mtr/binlog_streaming/t/pull_mode.test index 75201af..e5c2ccb 100644 --- a/mtr/binlog_streaming/t/pull_mode.test +++ b/mtr/binlog_streaming/t/pull_mode.test @@ -1,7 +1,3 @@ -# temporarily marking this test as non-parallel because of the hardcoded -# listening port ---source include/not_parallel.inc - --source ../include/have_binsrv.inc --source ../include/v80_v84_compatibility_defines.inc From 468826f141fcc37a6fc819746fcf7a60ed13f62f Mon Sep 17 00:00:00 2001 From: Kamil Holubicki Date: Mon, 21 Sep 2026 21:46:11 +0200 Subject: [PATCH 3/3] PBS-3 test: raise pull_mode grep-loop budget for ASAN + parallel pull_mode intentionally runs binsrv with $binsrv_read_timeout=3s and $binsrv_idle_time=1s so the source-mysqld restarts in the test force binsrv through several reconnect cycles. Each cycle is ~5s wall time and the test walks through three binlog files before checking that the fourth one appears. On release/debug builds the 60-iteration (60s) grep loop finishes with plenty of headroom, but under sanitized (ASAN) builds and parallel MTR execution each cycle stretches and 60s runs out before the fourth binlog shows up in the log. Bump $max_number_of_attempts from 60 to 300 (5 minutes). The loop exits on the first hit, so healthy runs pay nothing; it is still an order of magnitude smaller than the enclosing --testcase-timeout so a genuinely stuck binsrv still fails inside its intended envelope. Co-Authored-By: Claude Opus 4.7 --- mtr/binlog_streaming/t/pull_mode.test | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/mtr/binlog_streaming/t/pull_mode.test b/mtr/binlog_streaming/t/pull_mode.test index e5c2ccb..c89a0ea 100644 --- a/mtr/binlog_streaming/t/pull_mode.test +++ b/mtr/binlog_streaming/t/pull_mode.test @@ -97,7 +97,15 @@ FLUSH BINARY LOGS; --echo *** binary log. # We grep the Binlog Server Utility log file in a loop until we encounter the # fourth binary log file name. ---let $max_number_of_attempts = 60 +# The wait budget has to cover several reconnect cycles: the outbound +# read_timeout is 3s and the idle wait between reconnects is 1s (see the +# $binsrv_read_timeout / $binsrv_idle_time above), so binsrv needs +# multiple ~5s round trips to walk through the three preceding binlog +# files. Under sanitized (ASAN) builds and parallel MTR execution those +# round trips slow down further, so 60 iterations is not enough - 300 +# gives it several minutes and still bails out well before the outer +# testcase timeout kicks in. +--let $max_number_of_attempts = 300 --let $iteration = 0 while($iteration < $max_number_of_attempts) {