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/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..c89a0ea 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 @@ -101,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) { 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">()};