MySQL 26.7.0
Source Code Documentation
event_reader_controller.h
Go to the documentation of this file.
1// Copyright (c) 2026, Oracle and/or its affiliates.
2//
3// This program is free software; you can redistribute it and/or modify
4// it under the terms of the GNU General Public License, version 2.0,
5// as published by the Free Software Foundation.
6//
7// This program is designed to work with certain software (including
8// but not limited to OpenSSL) that is licensed under separate terms,
9// as designated in a particular file or component or in included license
10// documentation. The authors of MySQL hereby grant you an additional
11// permission to link the program and your derivative works with the
12// separately licensed software that they have either included with
13// the program or referenced in the documentation.
14//
15// This program is distributed in the hope that it will be useful,
16// but WITHOUT ANY WARRANTY; without even the implied warranty of
17// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
18// GNU General Public License, version 2.0, for more details.
19//
20// You should have received a copy of the GNU General Public License
21// along with this program; if not, write to the Free Software
22// Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA.
23
24#ifndef MYSQL_CSA_RELAY_LOG_EVENT_READER_CONTROLLER_H
25#define MYSQL_CSA_RELAY_LOG_EVENT_READER_CONTROLLER_H
26
27#include <memory>
28#include <optional>
44#include "sql/log_event.h" // Format_description_log_event
46#include "sql/rpl_rli.h" // Relay_log_info
47
48namespace mysql::csa {
49
50/// @brief The Event Reader / Controller class. This class uses the low level
51/// reader to read consecutive events
52/// from the relay log. Also, it exposed functions to remove
53/// consumed relay log files.
54/// Class provides methods to:
55/// - initialize internal data (open)
56/// - deinitialize internal data (close)
57/// - fetch the next event from the relay log (read_next)
58/// - fetch the next event metadata from the relay log
59/// (read_next)
60/// - register consumed relay log for purge and purge registered consecutive
61/// logs (concurrent_purge), which implements the `Log_purge_controller`
62/// interface
63/// @details This class may be seen as 'Rpl_applier_reader' created for the CSA.
64/// `Relay_log_decoder` works with prefetched relay log files or stream build
65/// on top of the IO_CACHE. The first one is used when reading from inactive
66/// files. When reading from active files, `Event_reader_controller` utilizes
67/// the IO_CACHE implementation, since it allows the applier to fetch
68/// data from cache instead of fetching data from disk. Data will be consumed
69/// from cache provided that the applier keeps up with Receiver thread and
70/// cache is 'large' enough to keep the recent data.
71/// When reading from inactive files, the 'Event_reader_controller' utilizes
72/// the prefetcher class utility to fetch data. When reading from active
73/// files, the Event Reader / Controller needs to rely on the 'MYSQL_BIN_LOG'
74/// synchronization primitives and supply the implementation of passive
75/// waiting for data (is_data_available, wait_data_ready).
76/// Following the legacy design, streams implemented on top of prefetcher
77/// are allowed to move to the next file upon the caller request. Therefore,
78/// the 'Event_reader_controller' is responsible for checking file boundaries
79/// and reopening the streams on top of new files when needed.
81 public:
83 /// Opens the first relay log
84 /// @retval false Success
85 /// @retval true failure
86 bool open();
87 /// Closes readers, stops prefetcher, clears internal state including error
88 /// state
89 void close();
90
91 // template <Reader_controller_read_type read_type>
92 // Reader_return_type<read_type>::type read(unsigned int return_timeout_ms) {
93 // if constexpr (read_type == Reader_controller_read_type::event) {
94 // auto tt =
95 // } else if constexpr (read_type == Reader_controller_read_type::metadata)
96 // { } else if constexpr (read_type ==
97 // Reader_controller_read_type::metadata) { } else {
98 // static_assert("not supported");
99 // }
100 // }
101
102 /// @brief Fetches next: event data, event metadata, event metadata
103 /// plus raw payload, depending on read type
104 /// @param return_timeout_ms If specified, will wait up to `return_timeout_ms`
105 /// miliseconds. If timeout is reached, returns true.
106 /// @param read_type Type of read: full event, event metadata or event
107 /// metadata and raw payload
108 /// @return Next: event data, event metadata, event metadata + raw payload
109 /// depending on read type. Empty object if an error occurred
110 std::optional<Event_file_metadata> read_next(
111 unsigned int return_timeout_ms, Reader_controller_read_type read_type);
112
113 /// @brief Register log_filename as a log ready to be purged. This function
114 /// will purge registered logs if they are in order according to
115 /// the index file content.
116 /// @param[in] log_filename Registeres this log as a log ready to be purged.
117 /// If registered logs for purging are in order, purges up to log_filename,
118 /// included.
119 bool concurrent_purge(const std::string &log_filename) override;
120
121 /// @brief Obtain currently processed file
122 /// @return Current file name
123 const std::string &get_file_name() const { return m_file_name; }
124
125 /// @brief Check if reader has an error
126 /// @return True in case an error occurred, false otherwise
127 bool is_error() const { return m_is_error; }
128
129 /// Stop the reader
130 void stop();
131
132 /// Check if reader is stopped
133 /// @return True when stop was requested; false otherwise
134 bool is_stopped() const;
135
136 private:
137 /// When next_log is true, opens the next file. Otherwise, opens the current
138 /// file
139 /// @param next_log When true, moves to next log after the current
140 /// @param offset Requested file offset
141 /// @retval false Success
142 /// @retval true failure
143 bool move_to_log(bool next_log = true, my_off_t offset = 0);
144
145 /// Moves to the next log file
146 /// @retval false Success
147 /// @retval true failure
149
150 /// Purge relay log files up to to_log
151 /// @param to_log Purge logs up to this log, exclusively
152 /// @retval false Success
153 /// @retval true Error
154 bool purge_applied_logs(const char *to_log);
155
156 /// In case we read from active file, we use this function to passively
157 /// wait for new event
158 /// @param return_timeout_ms Timeout after which we will exit from waiting
159 /// @retval false New event is available
160 /// @retval true Timeout occurred
161 bool wait_for_new_event(unsigned int return_timeout_ms);
162
163 /// @brief If cache is truncated, reopen the reader to avoid reading trash
164 /// data
165 /// @details Hack function that solves problem of relay log IO_CACHE
166 /// truncation on active relay log files
167 /// @return true if failure when reopening the file, false on success
168 /// @see Rpl_applier_reader::reopen_log_reader_if_needed
170
171 /// Sets internal error to msg
172 /// @param msg Error message
173 /// @return Error state : true
174 bool set_error(const char *msg);
175
176 /// Checks whether there is data in the current file
177 /// @retval true Data is available
178 /// @retval false No data
179 bool is_data_available();
180
181 /// Waits until data is available, stopped or return_timeout_ms is reached
182 /// @param return_timeout_ms Timeout after which we will exit from waiting
183 /// @retval true Data is available
184 /// @retval false Data is not ready / stopped / error
185 bool wait_data_ready(unsigned int return_timeout_ms);
186
187 /// Implements internal logic to choose between active and inactive file
188 /// reading
189 /// @param next_log True if open was called on the next log
190 /// @param prev_file Previous relay log file processed
191 void choose_reader(bool next_log, const std::string &prev_file);
192
193 /// Flag which is true in case any error occurred. Otherwise, it is set to
194 /// false.
195 bool m_is_error{false};
196 /// Error message if any
197 std::string m_error_msg{""};
198 /// non-owning RLI object pointer
200 /// Relay log prefetcher
202 /// Reader for active files
204 /// Reader for inactive files
206 /// Non-owning pointer to currently used reader (m_inactive_reader or
207 /// m_active_reader)
209 /// @brief Stores the current file as obtaining from stream
210 std::string m_file_name;
211 /// Here we keep the list of logs to be purged (in case later logs are
212 /// applied before), protected with m_rli->data_lock
213 std::unordered_set<std::string> m_logs_to_purge;
214 /// Flag indicating whether currently opened log file is active
216 /// This flag is true in case we are using active file reader and
217 /// reading from the active relay log file
219 /// Variable to decide on whether we want to run prefetcher
221 /// Stop flag
222 std::atomic<bool> m_is_stopped{false};
223};
224
225} // namespace mysql::csa
226
227#endif // MYSQL_CSA_RELAY_LOG_EVENT_READER_CONTROLLER_H
Interface class that all specializations of template <...> Basic_binlog_file_reader inherit from.
Definition: binlog_reader.h:364
Definition: rpl_rli.h:208
The Event Reader / Controller class.
Definition: event_reader_controller.h:80
bool is_error() const
Check if reader has an error.
Definition: event_reader_controller.h:127
Relay_log_info * m_rli
non-owning RLI object pointer
Definition: event_reader_controller.h:199
void close()
Closes readers, stops prefetcher, clears internal state including error state.
Definition: event_reader_controller.cpp:131
bool m_enable_prefetcher
Variable to decide on whether we want to run prefetcher.
Definition: event_reader_controller.h:220
bool m_using_prefetcher
Flag indicating whether currently opened log file is active.
Definition: event_reader_controller.h:215
void stop()
Stop the reader.
Definition: event_reader_controller.cpp:234
bool check_cache_truncated()
If cache is truncated, reopen the reader to avoid reading trash data.
Definition: event_reader_controller.cpp:142
Relaylog_file_reader m_active_reader
Reader for active files.
Definition: event_reader_controller.h:203
const std::string & get_file_name() const
Obtain currently processed file.
Definition: event_reader_controller.h:123
bool m_active_file_reading
This flag is true in case we are using active file reader and reading from the active relay log file.
Definition: event_reader_controller.h:218
std::string m_error_msg
Error message if any.
Definition: event_reader_controller.h:197
std::unordered_set< std::string > m_logs_to_purge
Here we keep the list of logs to be purged (in case later logs are applied before),...
Definition: event_reader_controller.h:213
IBasic_binlog_file_reader * m_current_reader
Non-owning pointer to currently used reader (m_inactive_reader or m_active_reader)
Definition: event_reader_controller.h:208
bool purge_applied_logs(const char *to_log)
Purge relay log files up to to_log.
Definition: event_reader_controller.cpp:354
Event_reader_controller(Relay_log_info *rli, Log_prefetcher_sptr prefetcher)
Definition: event_reader_controller.cpp:34
std::string m_file_name
Stores the current file as obtaining from stream.
Definition: event_reader_controller.h:210
bool move_to_next_log()
Moves to the next log file.
Log_prefetcher_sptr m_prefetcher
Relay log prefetcher.
Definition: event_reader_controller.h:201
bool set_error(const char *msg)
Sets internal error to msg.
Definition: event_reader_controller.cpp:53
bool is_stopped() const
Check if reader is stopped.
Definition: event_reader_controller.cpp:230
bool is_data_available()
Checks whether there is data in the current file.
Definition: event_reader_controller.cpp:155
bool wait_data_ready(unsigned int return_timeout_ms)
Waits until data is available, stopped or return_timeout_ms is reached.
Definition: event_reader_controller.cpp:174
bool concurrent_purge(const std::string &log_filename) override
Register log_filename as a log ready to be purged.
Definition: event_reader_controller.cpp:400
std::optional< Event_file_metadata > read_next(unsigned int return_timeout_ms, Reader_controller_read_type read_type)
Fetches next: event data, event metadata, event metadata plus raw payload, depending on read type.
Definition: event_reader_controller.cpp:242
bool move_to_log(bool next_log=true, my_off_t offset=0)
When next_log is true, opens the next file.
Definition: event_reader_controller.cpp:97
bool open()
Opens the first relay log.
Definition: event_reader_controller.cpp:64
bool m_is_error
Flag which is true in case any error occurred.
Definition: event_reader_controller.h:195
std::atomic< bool > m_is_stopped
Stop flag.
Definition: event_reader_controller.h:222
void choose_reader(bool next_log, const std::string &prev_file)
Implements internal logic to choose between active and inactive file reading.
Definition: event_reader_controller.cpp:79
Prefetched_relaylog_reader m_inactive_reader
Reader for inactive files.
Definition: event_reader_controller.h:205
bool wait_for_new_event(unsigned int return_timeout_ms)
In case we read from active file, we use this function to passively wait for new event.
Interface for a relay log controller which is capable of purging the relay logs.
Definition: log_purge_controller.h:39
Binary log event definitions.
ulonglong my_off_t
Definition: my_inttypes.h:72
static const char * log_filename
Definition: myisamlog.cc:97
constexpr bool prefetcher_enable
Definition: tune.h:41
Definition: channel.cpp:28
std::shared_ptr< Log_prefetcher > Log_prefetcher_sptr
Definition: log_prefetcher.h:48
Reader_controller_read_type
Definition: reader_controller_read_type.h:31
Experimental API header.
Transaction boundary parser definitions.