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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
80 changes: 80 additions & 0 deletions mtr/binlog_streaming/r/rs_blocking_dump.result
Original file line number Diff line number Diff line change
@@ -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 <BINSRV_STORAGE_PATH> 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.
153 changes: 153 additions & 0 deletions mtr/binlog_streaming/t/rs_blocking_dump.test
Original file line number Diff line number Diff line change
@@ -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
95 changes: 62 additions & 33 deletions src/minimysql/network_service.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@
#pragma GCC diagnostic pop

#include <boost/asio/redirect_error.hpp>
#include <boost/asio/steady_timer.hpp>
#include <boost/asio/this_coro.hpp>
#include <boost/asio/use_awaitable.hpp>

#include <boost/asio/ip/tcp.hpp>
Expand Down Expand Up @@ -274,6 +276,63 @@ void handle_exception(binsrv::basic_logger &logger, std::string_view context) {
}
}

[[nodiscard]] boost::asio::awaitable<void> 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"

Expand Down Expand Up @@ -517,39 +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};

// TODO: initialize sender_context with binlog_name:position extracted
// from the COM_BINLOG_DUMP command.
operations::sender_context sender_ctx{logger, storage, block_size};
bool fetch_result{};
util::const_byte_span event_data{};
while ((fetch_result = sender_ctx.get_event(event_data)) &&
!event_data.empty()) {
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 (!fetch_result) {
logger->log_format(binsrv::log_severity::error,
"net : failed to fetch next event block for {}",
remote_endpoint_str);
terminated = true;
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: {
Expand Down
10 changes: 5 additions & 5 deletions src/operations/sender_context.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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{};
Expand Down
Loading