Skip to main content

mysql_common/binlog/events/
mod.rs

1// Copyright (c) 2021 Anatoly Ikorsky
2//
3// Licensed under the Apache License, Version 2.0
4// <LICENSE-APACHE or http://www.apache.org/licenses/LICENSE-2.0> or the MIT
5// license <LICENSE-MIT or http://opensource.org/licenses/MIT>, at your
6// option. All files in the project carrying such notice may not be copied,
7// modified, or distributed except according to those terms.
8
9use bitvec::prelude::*;
10use byteorder::{LittleEndian, WriteBytesExt};
11use bytes::BufMut;
12use saturating::Saturating as S;
13
14use crate::{
15    io::ParseBuf,
16    misc::raw::{RawConst, RawFlags, int::*},
17    proto::{MyDeserialize, MySerialize},
18};
19
20pub use self::{
21    anonymous_gtid_event::AnonymousGtidEvent,
22    begin_load_query_event::BeginLoadQueryEvent,
23    delete_rows_event::DeleteRowsEvent,
24    delete_rows_event_v1::DeleteRowsEventV1,
25    execute_load_query_event::ExecuteLoadQueryEvent,
26    format_description_event::FormatDescriptionEvent,
27    gtid_event::GtidEvent,
28    incident_event::IncidentEvent,
29    intvar_event::IntvarEvent,
30    partial_update_rows_event::PartialUpdateRowsEvent,
31    previous_gtids_event::{PreviousGtidsEvent, PreviousGtidsSid},
32    query_event::{QueryEvent, StatusVar, StatusVarVal, StatusVars, StatusVarsIterator},
33    rand_event::RandEvent,
34    rotate_event::RotateEvent,
35    rows_event::{RowsEvent, RowsEventRows},
36    rows_query_event::RowsQueryEvent,
37    table_map_event::*,
38    transaction_payload_event::{TransactionPayloadEvent, TransactionPayloadReader},
39    update_rows_event::UpdateRowsEvent,
40    update_rows_event_v1::UpdateRowsEventV1,
41    user_var_event::UserVarEvent,
42    write_rows_event::WriteRowsEvent,
43    write_rows_event_v1::WriteRowsEventV1,
44    xid_event::XidEvent,
45};
46
47use std::{
48    any::type_name,
49    borrow::Cow,
50    cmp::min,
51    io::{self, Read, Write},
52};
53
54use super::{
55    BinlogCtx, BinlogEvent,
56    consts::{
57        BinlogChecksumAlg, BinlogVersion, EventFlags, EventType, RowsEventFlags,
58        UnknownChecksumAlg, UnknownEventType,
59    },
60    misc::LimitWrite,
61};
62
63mod anonymous_gtid_event;
64mod begin_load_query_event;
65mod delete_rows_event;
66mod delete_rows_event_v1;
67mod execute_load_query_event;
68mod format_description_event;
69mod gtid_event;
70
71mod incident_event;
72mod intvar_event;
73mod partial_update_rows_event;
74mod previous_gtids_event;
75mod query_event;
76mod rand_event;
77mod rotate_event;
78mod rows_event;
79mod rows_query_event;
80mod table_map_event;
81mod transaction_payload_event;
82mod update_rows_event;
83mod update_rows_event_v1;
84mod user_var_event;
85mod write_rows_event;
86mod write_rows_event_v1;
87mod xid_event;
88
89/// Raw binlog event.
90///
91/// A binlog event starts with a Binlog Event header and is followed by a Binlog Event Type
92/// specific data part.
93#[derive(Debug, Clone, Eq, PartialEq)]
94pub struct Event {
95    /// Format description event.
96    fde: FormatDescriptionEvent<'static>,
97    /// Common header of an event.
98    header: BinlogEventHeader,
99    /// An event-type specific data.
100    ///
101    /// Checksum-related suffix is truncated:
102    ///
103    /// *   checksum algorithm description (for fde) will go to `footer`;
104    /// *   checksum will go to `checksum`.
105    data: Vec<u8>,
106    /// Log event footer.
107    footer: BinlogEventFooter,
108    /// Event checksum.
109    ///
110    /// Makes sense only if checksum algorithm is defined in `footer`.
111    checksum: [u8; BinlogEventFooter::BINLOG_CHECKSUM_LEN],
112}
113
114impl Event {
115    /// Reads an event from `input`.
116    pub fn read<'a, T: Read>(
117        fde: &'a FormatDescriptionEvent<'a>,
118        mut input: T,
119    ) -> io::Result<Self> {
120        let binlog_header_len = BinlogEventHeader::LEN;
121        let mut fde = fde.clone().into_owned();
122
123        let mut header_buf = [0u8; BinlogEventHeader::LEN];
124        input.read_exact(&mut header_buf)?;
125        let header = BinlogEventHeader::deserialize((), &mut ParseBuf(&header_buf))?;
126
127        let mut data = vec![0_u8; (S(header.event_size() as usize) - S(binlog_header_len)).0];
128        input.read_exact(&mut data).unwrap();
129
130        let is_fde = header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
131        let mut bytes_to_truncate = 0;
132        let mut checksum = [0_u8; BinlogEventFooter::BINLOG_CHECKSUM_LEN];
133
134        let footer = if is_fde {
135            let footer = BinlogEventFooter::read(&data)?;
136            if footer.checksum_alg.is_some() {
137                // truncate checksum algorithm description
138                bytes_to_truncate += BinlogEventFooter::BINLOG_CHECKSUM_ALG_DESC_LEN;
139            }
140            // We'll update dummy fde footer
141            fde = fde.with_footer(footer);
142            footer
143        } else {
144            fde.footer()
145        };
146
147        // * fde will always contain checksum (see WL#2540)
148        // * events inside of a Transaction_payload_event are not checksummed (see WL#3549)
149        let contains_checksum = footer.checksum_alg.is_some()
150            && (is_fde || footer.checksum_alg != Some(RawConst::new(0)))
151            && footer.checksum_enabled;
152
153        if contains_checksum {
154            // truncate checksum
155            bytes_to_truncate += BinlogEventFooter::BINLOG_CHECKSUM_LEN;
156            checksum.copy_from_slice(&data[data.len() - BinlogEventFooter::BINLOG_CHECKSUM_LEN..]);
157        }
158
159        data.truncate(data.len() - bytes_to_truncate);
160
161        Ok(Self {
162            fde,
163            header,
164            data,
165            footer,
166            checksum,
167        })
168    }
169
170    /// Writes this event into the `output`.
171    pub fn write<T: Write>(&self, version: BinlogVersion, mut output: T) -> io::Result<()> {
172        let is_fde = self.header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
173        let mut output = output.limit(S(self.len(version)));
174
175        let mut header_buf = Vec::with_capacity(BinlogEventHeader::LEN);
176        self.header.serialize(&mut header_buf);
177        output.write_all(&header_buf)?;
178        output.write_all(&self.data)?;
179
180        if let Ok(Some(alg)) = self.footer.get_checksum_alg() {
181            if is_fde {
182                output.write_u8(alg as u8)?;
183            }
184            if alg == BinlogChecksumAlg::BINLOG_CHECKSUM_ALG_CRC32 || is_fde {
185                output.write_u32::<LittleEndian>(self.calc_checksum(alg))?;
186            }
187        }
188
189        Ok(())
190    }
191
192    /// Returns a length of a serialized representation of this event.
193    fn len(&self, _version: BinlogVersion) -> usize {
194        let is_fde = self.header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
195        let mut len = S(0);
196
197        len += S(BinlogEventHeader::LEN);
198        len += S(self.data.len());
199        if let Ok(Some(alg)) = self.footer.get_checksum_alg() {
200            if is_fde {
201                len += S(BinlogEventFooter::BINLOG_CHECKSUM_ALG_DESC_LEN);
202            }
203            if is_fde || alg != BinlogChecksumAlg::BINLOG_CHECKSUM_ALG_OFF {
204                len += S(BinlogEventFooter::BINLOG_CHECKSUM_LEN);
205            }
206        }
207
208        min(len.0, u32::MAX as usize - BinlogEventHeader::LEN)
209    }
210
211    /// Returns a reference to the corresponding format description event.
212    pub fn fde(&self) -> &FormatDescriptionEvent<'static> {
213        &self.fde
214    }
215
216    /// Returns a reference to the event header.
217    pub fn header(&self) -> BinlogEventHeader {
218        self.header
219    }
220
221    /// Returns a reference to the event data.
222    pub fn data(&self) -> &[u8] {
223        &self.data
224    }
225
226    /// Returns a reference to the event footer.
227    pub fn footer(&self) -> BinlogEventFooter {
228        self.footer
229    }
230
231    /// Returns the checksum, if it is defined.
232    pub fn checksum(&self) -> Option<[u8; BinlogEventFooter::BINLOG_CHECKSUM_LEN]> {
233        let contains_checksum = self.footer.checksum_alg.is_some()
234            && (self.header.event_type.0 == (EventType::FORMAT_DESCRIPTION_EVENT as u8)
235                || self.footer.checksum_alg != Some(RawConst::new(0)));
236        contains_checksum.then_some(self.checksum)
237    }
238
239    /// Read event-type specific data as a binlog struct.
240    pub fn read_event<'a, T: BinlogEvent<'a>>(&'a self) -> io::Result<T> {
241        // we'll use data.len() here because of truncated event footer
242        let event_size = BinlogEventHeader::LEN + self.data.len();
243        let event_data = &mut ParseBuf(&self.data);
244        let ctx = BinlogCtx::new(event_size, &self.fde, self.header.event_type_raw());
245
246        let event = event_data.parse(ctx)?;
247
248        // it is an error if the `event_data` isn't fully consumed
249        if !event_data.is_empty() {
250            return Err(io::Error::other(format!(
251                "bytes remaining on stream while reading {}",
252                type_name::<T>()
253            )));
254        }
255
256        Ok(event)
257    }
258
259    /// Reads event data. Returns `None` if event type is unknown.
260    pub fn read_data(&self) -> io::Result<Option<EventData<'_>>> {
261        use EventType::*;
262
263        let event_type = match self.header.event_type.get() {
264            Ok(event_type) => event_type,
265            _ => return Ok(None),
266        };
267
268        let event_data = match event_type {
269            ENUM_END_EVENT | UNKNOWN_EVENT => EventData::UnknownEvent,
270            START_EVENT_V3 => EventData::StartEventV3(Cow::Borrowed(&*self.data)),
271            QUERY_EVENT => EventData::QueryEvent(self.read_event()?),
272            STOP_EVENT => EventData::StopEvent,
273            ROTATE_EVENT => EventData::RotateEvent(self.read_event()?),
274            INTVAR_EVENT => EventData::IntvarEvent(self.read_event()?),
275            LOAD_EVENT => EventData::LoadEvent(Cow::Borrowed(&*self.data)),
276            SLAVE_EVENT => EventData::SlaveEvent,
277            CREATE_FILE_EVENT => EventData::CreateFileEvent(Cow::Borrowed(&*self.data)),
278            APPEND_BLOCK_EVENT => EventData::AppendBlockEvent(Cow::Borrowed(&*self.data)),
279            EXEC_LOAD_EVENT => EventData::ExecLoadEvent(Cow::Borrowed(&*self.data)),
280            DELETE_FILE_EVENT => EventData::DeleteFileEvent(Cow::Borrowed(&*self.data)),
281            NEW_LOAD_EVENT => EventData::NewLoadEvent(Cow::Borrowed(&*self.data)),
282            RAND_EVENT => EventData::RandEvent(self.read_event()?),
283            USER_VAR_EVENT => EventData::UserVarEvent(self.read_event()?),
284            FORMAT_DESCRIPTION_EVENT => {
285                let fde = self
286                    .read_event::<FormatDescriptionEvent<'_>>()?
287                    .with_footer(self.footer);
288                EventData::FormatDescriptionEvent(fde)
289            }
290            XID_EVENT => EventData::XidEvent(self.read_event()?),
291            BEGIN_LOAD_QUERY_EVENT => EventData::BeginLoadQueryEvent(self.read_event()?),
292            EXECUTE_LOAD_QUERY_EVENT => EventData::ExecuteLoadQueryEvent(self.read_event()?),
293            TABLE_MAP_EVENT => EventData::TableMapEvent(self.read_event()?),
294            PRE_GA_WRITE_ROWS_EVENT => EventData::PreGaWriteRowsEvent(Cow::Borrowed(&*self.data)),
295            PRE_GA_UPDATE_ROWS_EVENT => EventData::PreGaUpdateRowsEvent(Cow::Borrowed(&*self.data)),
296            PRE_GA_DELETE_ROWS_EVENT => EventData::PreGaDeleteRowsEvent(Cow::Borrowed(&*self.data)),
297            WRITE_ROWS_EVENT_V1 => {
298                EventData::RowsEvent(RowsEventData::WriteRowsEventV1(self.read_event()?))
299            }
300            UPDATE_ROWS_EVENT_V1 => {
301                EventData::RowsEvent(RowsEventData::UpdateRowsEventV1(self.read_event()?))
302            }
303            DELETE_ROWS_EVENT_V1 => {
304                EventData::RowsEvent(RowsEventData::DeleteRowsEventV1(self.read_event()?))
305            }
306            INCIDENT_EVENT => EventData::IncidentEvent(self.read_event()?),
307            HEARTBEAT_EVENT => EventData::HeartbeatEvent,
308            IGNORABLE_EVENT => EventData::IgnorableEvent(Cow::Borrowed(&*self.data)),
309            ROWS_QUERY_EVENT => EventData::RowsQueryEvent(self.read_event()?),
310            WRITE_ROWS_EVENT => {
311                EventData::RowsEvent(RowsEventData::WriteRowsEvent(self.read_event()?))
312            }
313            UPDATE_ROWS_EVENT => {
314                EventData::RowsEvent(RowsEventData::UpdateRowsEvent(self.read_event()?))
315            }
316            DELETE_ROWS_EVENT => {
317                EventData::RowsEvent(RowsEventData::DeleteRowsEvent(self.read_event()?))
318            }
319            GTID_EVENT => EventData::GtidEvent(self.read_event()?),
320            ANONYMOUS_GTID_EVENT => EventData::AnonymousGtidEvent(self.read_event()?),
321            PREVIOUS_GTIDS_EVENT => EventData::PreviousGtidsEvent(self.read_event()?),
322            TRANSACTION_CONTEXT_EVENT => {
323                EventData::TransactionContextEvent(Cow::Borrowed(&*self.data))
324            }
325            VIEW_CHANGE_EVENT => EventData::ViewChangeEvent(Cow::Borrowed(&*self.data)),
326            XA_PREPARE_LOG_EVENT => EventData::XaPrepareLogEvent(Cow::Borrowed(&*self.data)),
327            PARTIAL_UPDATE_ROWS_EVENT => {
328                EventData::RowsEvent(RowsEventData::PartialUpdateRowsEvent(self.read_event()?))
329            }
330            TRANSACTION_PAYLOAD_EVENT => EventData::TransactionPayloadEvent(self.read_event()?),
331            GTID_TAGGED_LOG_EVENT => EventData::GtidEvent(self.read_event()?),
332        };
333
334        Ok(Some(event_data))
335    }
336
337    /// Calculates checksum for this event.
338    pub fn calc_checksum(&self, alg: BinlogChecksumAlg) -> u32 {
339        let is_fde = self.header.event_type.0 == EventType::FORMAT_DESCRIPTION_EVENT as u8;
340
341        let mut hasher = crc32fast::Hasher::new();
342        let mut header = Vec::with_capacity(BinlogEventHeader::LEN);
343        let mut header_struct = self.header;
344        if header_struct
345            .flags
346            .get()
347            .contains(EventFlags::LOG_EVENT_BINLOG_IN_USE_F)
348        {
349            // In case this is a Format_description_log_event, we need to clear
350            // the LOG_EVENT_BINLOG_IN_USE_F flag before computing the checksum,
351            // since the flag will be cleared when the binlog is closed.
352            // On verification, the flag is also dropped before computing the checksum.
353            header_struct.flags.0 &= !(EventFlags::LOG_EVENT_BINLOG_IN_USE_F.bits());
354        }
355        header_struct.serialize(&mut header);
356        hasher.update(&header);
357        hasher.update(&self.data);
358        if is_fde {
359            hasher.update(&[alg as u8][..]);
360        }
361        hasher.finalize()
362    }
363}
364
365/// The binlog event header starts each event and is 19 bytes long assuming binlog version >= 4.
366#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
367pub struct BinlogEventHeader {
368    /// Seconds since unix epoch.
369    timestamp: RawInt<LeU32>,
370    /// Raw event Type.
371    event_type: RawConst<u8, EventType>,
372    /// Server-id of the originating mysql-server.
373    ///
374    /// Used to filter out events in circular replication.
375    server_id: RawInt<LeU32>,
376    /// Size of the event (header, post-header, body).
377    event_size: RawInt<LeU32>,
378    /// Position of the next event.
379    log_pos: RawInt<LeU32>,
380    /// Binlog Event Flag.
381    ///
382    /// This field contains the raw value. Use [`Self::flags()`] to get the actual flags.
383    flags: RawFlags<EventFlags, LeU16>,
384}
385
386impl BinlogEventHeader {
387    /// Binlog event header length for version >= 4.
388    pub const LEN: usize = 19;
389
390    /// Creates a new `BinlogEventHeader`.
391    pub fn new(
392        timestamp: u32,
393        event_type: EventType,
394        server_id: u32,
395        event_size: u32,
396        log_pos: u32,
397        flags: EventFlags,
398    ) -> Self {
399        Self {
400            timestamp: RawInt::new(timestamp),
401            event_type: RawConst::new(event_type as u8),
402            server_id: RawInt::new(server_id),
403            event_size: RawInt::new(event_size),
404            log_pos: RawInt::new(log_pos),
405            flags: RawFlags::new(flags.bits()),
406        }
407    }
408
409    /// Returns the `timestamp` value.
410    ///
411    /// `timestamp` is in seconds since unix epoch.
412    pub fn timestamp(&self) -> u32 {
413        self.timestamp.0
414    }
415
416    /// Returns the raw event type.
417    pub fn event_type_raw(&self) -> u8 {
418        self.event_type.0
419    }
420
421    /// Returns the event type, if it's valid.
422    pub fn event_type(&self) -> Result<EventType, UnknownEventType> {
423        self.event_type.get()
424    }
425
426    /// Returns the server Id of the originating mysql-server.
427    ///
428    /// Used to filter out events in circular replication.
429    pub fn server_id(&self) -> u32 {
430        self.server_id.0
431    }
432
433    /// Returns the size of the event (header, post-header, body).
434    pub fn event_size(&self) -> u32 {
435        self.event_size.0
436    }
437
438    /// Returns the position of the next event.
439    pub fn log_pos(&self) -> u32 {
440        self.log_pos.0
441    }
442
443    /// Returns event flags (unknown bits are truncated).
444    pub fn flags(&self) -> EventFlags {
445        self.flags.get()
446    }
447
448    /// Returns raw event flags (unknown bits are preserved).
449    pub fn flags_raw(&self) -> u16 {
450        self.flags.0
451    }
452}
453
454impl<'de> MyDeserialize<'de> for BinlogEventHeader {
455    const SIZE: Option<usize> = Some(Self::LEN);
456    type Ctx = ();
457
458    fn deserialize((): Self::Ctx, buf: &mut ParseBuf<'de>) -> io::Result<Self> {
459        let mut buf: ParseBuf<'_> = buf.parse_unchecked(Self::LEN)?;
460        Ok(Self {
461            timestamp: buf.parse_unchecked(())?,
462            event_type: buf.parse_unchecked(())?,
463            server_id: buf.parse_unchecked(())?,
464            event_size: buf.parse_unchecked(())?,
465            log_pos: buf.parse_unchecked(())?,
466            flags: buf.parse_unchecked(())?,
467        })
468    }
469}
470
471impl MySerialize for BinlogEventHeader {
472    fn serialize(&self, buf: &mut Vec<u8>) {
473        self.timestamp.serialize(&mut *buf);
474        self.event_type.serialize(&mut *buf);
475        self.server_id.serialize(&mut *buf);
476        self.event_size.serialize(&mut *buf);
477        self.log_pos.serialize(&mut *buf);
478        self.flags.serialize(&mut *buf);
479    }
480}
481
482/// Binlog event footer.
483#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
484pub struct BinlogEventFooter {
485    /// Raw checksum algorithm description.
486    checksum_alg: Option<RawConst<u8, BinlogChecksumAlg>>,
487
488    /// Checksum enabled
489    checksum_enabled: bool,
490}
491
492impl BinlogEventFooter {
493    /// Length of the checksum algorithm description.
494    pub const BINLOG_CHECKSUM_ALG_DESC_LEN: usize = 1;
495    /// Length of the checksum.
496    pub const BINLOG_CHECKSUM_LEN: usize = 4;
497    /// Minimum MySql version that supports checksums.
498    pub const CHECKSUM_VERSION_PRODUCT: (u8, u8, u8) = (5, 6, 1);
499
500    pub fn new(checksum_alg: BinlogChecksumAlg) -> Self {
501        Self {
502            checksum_alg: Some(RawConst::new(checksum_alg as u8)),
503            checksum_enabled: true,
504        }
505    }
506
507    /// Returns parsed checksum algorithm, or raw value if algorithm is unknown.
508    pub fn get_checksum_alg(&self) -> Result<Option<BinlogChecksumAlg>, UnknownChecksumAlg> {
509        self.checksum_alg.as_ref().map(RawConst::get).transpose()
510    }
511
512    /// Returns `true` if checksum is enabled.
513    pub fn get_checksum_enabled(&self) -> bool {
514        self.checksum_enabled
515    }
516
517    /// Set checksum enabled flag.
518    pub fn set_checksum_enabled(&mut self, enabled: bool) {
519        self.checksum_enabled = enabled;
520    }
521
522    /// Reads binlog event footer from the given buffer.
523    ///
524    /// Requires that buf contains `FormatDescriptionEvent` data.
525    pub fn read(buf: &[u8]) -> io::Result<Self> {
526        let checksum_alg = if buf.len()
527            >= FormatDescriptionEvent::SERVER_VER_OFFSET + FormatDescriptionEvent::SERVER_VER_LEN
528        {
529            let mut server_version = vec![0_u8; FormatDescriptionEvent::SERVER_VER_LEN];
530            (&buf[FormatDescriptionEvent::SERVER_VER_OFFSET..]).read_exact(&mut server_version)?;
531            server_version[FormatDescriptionEvent::SERVER_VER_LEN - 1] = 0;
532            let version = crate::misc::split_version(&server_version);
533            if version < Self::CHECKSUM_VERSION_PRODUCT {
534                None
535            } else {
536                let offset = buf.len()
537                    - (BinlogEventFooter::BINLOG_CHECKSUM_ALG_DESC_LEN
538                        + BinlogEventFooter::BINLOG_CHECKSUM_LEN);
539                Some(buf[offset])
540            }
541        } else {
542            None
543        };
544
545        Ok(Self {
546            checksum_alg: checksum_alg.map(RawConst::new),
547            checksum_enabled: true,
548        })
549    }
550}
551
552impl Default for BinlogEventFooter {
553    fn default() -> Self {
554        BinlogEventFooter {
555            checksum_alg: Some(RawConst::new(
556                BinlogChecksumAlg::BINLOG_CHECKSUM_ALG_OFF as u8,
557            )),
558            checksum_enabled: true,
559        }
560    }
561}
562
563/// Parsed event data.
564#[derive(Debug, Clone, Eq, PartialEq, Hash)]
565pub enum EventData<'a> {
566    UnknownEvent,
567    /// Ignored by this implementation
568    StartEventV3(Cow<'a, [u8]>),
569    QueryEvent(QueryEvent<'a>),
570    StopEvent,
571    RotateEvent(RotateEvent<'a>),
572    IntvarEvent(IntvarEvent),
573    /// Ignored by this implementation
574    LoadEvent(Cow<'a, [u8]>),
575    SlaveEvent,
576    CreateFileEvent(Cow<'a, [u8]>),
577    /// Ignored by this implementation
578    AppendBlockEvent(Cow<'a, [u8]>),
579    /// Ignored by this implementation
580    ExecLoadEvent(Cow<'a, [u8]>),
581    /// Ignored by this implementation
582    DeleteFileEvent(Cow<'a, [u8]>),
583    /// Ignored by this implementation
584    NewLoadEvent(Cow<'a, [u8]>),
585    RandEvent(RandEvent),
586    UserVarEvent(UserVarEvent<'a>),
587    FormatDescriptionEvent(FormatDescriptionEvent<'a>),
588    XidEvent(XidEvent),
589    BeginLoadQueryEvent(BeginLoadQueryEvent<'a>),
590    ExecuteLoadQueryEvent(ExecuteLoadQueryEvent<'a>),
591    TableMapEvent(TableMapEvent<'a>),
592    /// Ignored by this implementation
593    PreGaWriteRowsEvent(Cow<'a, [u8]>),
594    /// Ignored by this implementation
595    PreGaUpdateRowsEvent(Cow<'a, [u8]>),
596    /// Ignored by this implementation
597    PreGaDeleteRowsEvent(Cow<'a, [u8]>),
598    IncidentEvent(IncidentEvent<'a>),
599    HeartbeatEvent,
600    IgnorableEvent(Cow<'a, [u8]>),
601    RowsQueryEvent(RowsQueryEvent<'a>),
602    GtidEvent(GtidEvent),
603    /// Not yet implemented.
604    AnonymousGtidEvent(AnonymousGtidEvent),
605    PreviousGtidsEvent(PreviousGtidsEvent<'a>),
606    /// Not yet implemented.
607    TransactionContextEvent(Cow<'a, [u8]>),
608    /// Not yet implemented.
609    ViewChangeEvent(Cow<'a, [u8]>),
610    /// Not yet implemented.
611    XaPrepareLogEvent(Cow<'a, [u8]>),
612    RowsEvent(RowsEventData<'a>),
613    TransactionPayloadEvent(TransactionPayloadEvent<'a>),
614}
615
616impl<'a> EventData<'a> {
617    pub fn into_owned(self) -> EventData<'static> {
618        match self {
619            EventData::UnknownEvent => EventData::UnknownEvent,
620            EventData::StartEventV3(ev) => EventData::StartEventV3(Cow::Owned(ev.into_owned())),
621            Self::QueryEvent(ev) => EventData::QueryEvent(ev.into_owned()),
622            Self::StopEvent => EventData::StopEvent,
623            Self::RotateEvent(ev) => EventData::RotateEvent(ev.into_owned()),
624            Self::IntvarEvent(ev) => EventData::IntvarEvent(ev),
625            Self::LoadEvent(ev) => EventData::LoadEvent(Cow::Owned(ev.into_owned())),
626            Self::SlaveEvent => EventData::SlaveEvent,
627            Self::CreateFileEvent(ev) => EventData::CreateFileEvent(Cow::Owned(ev.into_owned())),
628            Self::AppendBlockEvent(ev) => EventData::AppendBlockEvent(Cow::Owned(ev.into_owned())),
629            Self::ExecLoadEvent(ev) => EventData::ExecLoadEvent(Cow::Owned(ev.into_owned())),
630            Self::DeleteFileEvent(ev) => EventData::DeleteFileEvent(Cow::Owned(ev.into_owned())),
631            Self::NewLoadEvent(ev) => EventData::NewLoadEvent(Cow::Owned(ev.into_owned())),
632            Self::RandEvent(ev) => EventData::RandEvent(ev),
633            Self::UserVarEvent(ev) => EventData::UserVarEvent(ev.into_owned()),
634            Self::FormatDescriptionEvent(ev) => EventData::FormatDescriptionEvent(ev.into_owned()),
635            Self::XidEvent(ev) => EventData::XidEvent(ev),
636            Self::BeginLoadQueryEvent(ev) => EventData::BeginLoadQueryEvent(ev.into_owned()),
637            Self::ExecuteLoadQueryEvent(ev) => EventData::ExecuteLoadQueryEvent(ev.into_owned()),
638            Self::TableMapEvent(ev) => EventData::TableMapEvent(ev.into_owned()),
639            Self::PreGaWriteRowsEvent(ev) => {
640                EventData::PreGaWriteRowsEvent(Cow::Owned(ev.into_owned()))
641            }
642            Self::PreGaUpdateRowsEvent(ev) => {
643                EventData::PreGaUpdateRowsEvent(Cow::Owned(ev.into_owned()))
644            }
645            Self::PreGaDeleteRowsEvent(ev) => {
646                EventData::PreGaDeleteRowsEvent(Cow::Owned(ev.into_owned()))
647            }
648            Self::IncidentEvent(ev) => EventData::IncidentEvent(ev.into_owned()),
649            Self::HeartbeatEvent => EventData::HeartbeatEvent,
650            Self::IgnorableEvent(ev) => EventData::IgnorableEvent(Cow::Owned(ev.into_owned())),
651            Self::RowsQueryEvent(ev) => EventData::RowsQueryEvent(ev.into_owned()),
652            Self::GtidEvent(ev) => EventData::GtidEvent(ev),
653            Self::AnonymousGtidEvent(ev) => EventData::AnonymousGtidEvent(ev),
654            Self::PreviousGtidsEvent(ev) => EventData::PreviousGtidsEvent(ev.into_owned()),
655            Self::TransactionContextEvent(ev) => {
656                EventData::TransactionContextEvent(Cow::Owned(ev.into_owned()))
657            }
658            Self::ViewChangeEvent(ev) => EventData::ViewChangeEvent(Cow::Owned(ev.into_owned())),
659            Self::XaPrepareLogEvent(ev) => {
660                EventData::XaPrepareLogEvent(Cow::Owned(ev.into_owned()))
661            }
662            Self::RowsEvent(ev) => EventData::RowsEvent(ev.into_owned()),
663            Self::TransactionPayloadEvent(ev) => {
664                EventData::TransactionPayloadEvent(ev.into_owned())
665            }
666        }
667    }
668}
669
670impl MySerialize for EventData<'_> {
671    fn serialize(&self, buf: &mut Vec<u8>) {
672        match self {
673            EventData::UnknownEvent => (),
674            EventData::StartEventV3(ev) => buf.put_slice(ev),
675            EventData::QueryEvent(ev) => ev.serialize(buf),
676            EventData::StopEvent => (),
677            EventData::RotateEvent(ev) => ev.serialize(buf),
678            EventData::IntvarEvent(ev) => ev.serialize(buf),
679            EventData::LoadEvent(ev) => buf.put_slice(ev),
680            EventData::SlaveEvent => (),
681            EventData::CreateFileEvent(ev) => buf.put_slice(ev),
682            EventData::AppendBlockEvent(ev) => buf.put_slice(ev),
683            EventData::ExecLoadEvent(ev) => buf.put_slice(ev),
684            EventData::DeleteFileEvent(ev) => buf.put_slice(ev),
685            EventData::NewLoadEvent(ev) => buf.put_slice(ev),
686            EventData::RandEvent(ev) => ev.serialize(buf),
687            EventData::UserVarEvent(ev) => ev.serialize(buf),
688            EventData::FormatDescriptionEvent(ev) => ev.serialize(buf),
689            EventData::XidEvent(ev) => ev.serialize(buf),
690            EventData::BeginLoadQueryEvent(ev) => ev.serialize(buf),
691            EventData::ExecuteLoadQueryEvent(ev) => ev.serialize(buf),
692            EventData::TableMapEvent(ev) => ev.serialize(buf),
693            EventData::PreGaWriteRowsEvent(ev) => buf.put_slice(ev),
694            EventData::PreGaUpdateRowsEvent(ev) => buf.put_slice(ev),
695            EventData::PreGaDeleteRowsEvent(ev) => buf.put_slice(ev),
696            EventData::IncidentEvent(ev) => ev.serialize(buf),
697            EventData::HeartbeatEvent => (),
698            EventData::IgnorableEvent(ev) => buf.put_slice(ev),
699            EventData::RowsQueryEvent(ev) => ev.serialize(buf),
700            EventData::GtidEvent(ev) => ev.serialize(buf),
701            EventData::AnonymousGtidEvent(ev) => ev.serialize(buf),
702            EventData::PreviousGtidsEvent(ev) => ev.serialize(buf),
703            EventData::TransactionContextEvent(ev) => buf.put_slice(ev),
704            EventData::ViewChangeEvent(ev) => buf.put_slice(ev),
705            EventData::XaPrepareLogEvent(ev) => buf.put_slice(ev),
706            EventData::RowsEvent(ev) => ev.serialize(buf),
707            EventData::TransactionPayloadEvent(ev) => ev.serialize(buf),
708        }
709    }
710}
711
712/// Rows events are unified under this enum (see [`EventData`]).
713#[derive(Debug, Clone, Eq, PartialEq, Hash)]
714pub enum RowsEventData<'a> {
715    WriteRowsEventV1(WriteRowsEventV1<'a>),
716    UpdateRowsEventV1(UpdateRowsEventV1<'a>),
717    DeleteRowsEventV1(DeleteRowsEventV1<'a>),
718    WriteRowsEvent(WriteRowsEvent<'a>),
719    UpdateRowsEvent(UpdateRowsEvent<'a>),
720    DeleteRowsEvent(DeleteRowsEvent<'a>),
721    PartialUpdateRowsEvent(PartialUpdateRowsEvent<'a>),
722}
723
724impl<'a> RowsEventData<'a> {
725    /// Returns the number that identifies the table (see `TableMapEvent`).
726    pub fn table_id(&self) -> u64 {
727        match self {
728            RowsEventData::WriteRowsEventV1(ev) => ev.table_id(),
729            RowsEventData::UpdateRowsEventV1(ev) => ev.table_id(),
730            RowsEventData::DeleteRowsEventV1(ev) => ev.table_id(),
731            RowsEventData::WriteRowsEvent(ev) => ev.table_id(),
732            RowsEventData::UpdateRowsEvent(ev) => ev.table_id(),
733            RowsEventData::DeleteRowsEvent(ev) => ev.table_id(),
734            RowsEventData::PartialUpdateRowsEvent(ev) => ev.table_id(),
735        }
736    }
737
738    /// Returns the number of columns in the table.
739    pub fn num_columns(&self) -> u64 {
740        match self {
741            RowsEventData::WriteRowsEventV1(ev) => ev.num_columns(),
742            RowsEventData::UpdateRowsEventV1(ev) => ev.num_columns(),
743            RowsEventData::DeleteRowsEventV1(ev) => ev.num_columns(),
744            RowsEventData::WriteRowsEvent(ev) => ev.num_columns(),
745            RowsEventData::UpdateRowsEvent(ev) => ev.num_columns(),
746            RowsEventData::DeleteRowsEvent(ev) => ev.num_columns(),
747            RowsEventData::PartialUpdateRowsEvent(ev) => ev.num_columns(),
748        }
749    }
750
751    /// Returns columns in the before-image (only for DELETE and UPDATE).
752    ///
753    /// Each bit indicates whether corresponding column is used in the image.
754    pub fn columns_before_image(&'a self) -> Option<&'a BitSlice<u8>> {
755        match self {
756            RowsEventData::WriteRowsEventV1(_) => None,
757            RowsEventData::UpdateRowsEventV1(ev) => Some(ev.columns_before_image()),
758            RowsEventData::DeleteRowsEventV1(ev) => Some(ev.columns_before_image()),
759            RowsEventData::WriteRowsEvent(_) => None,
760            RowsEventData::UpdateRowsEvent(ev) => Some(ev.columns_before_image()),
761            RowsEventData::DeleteRowsEvent(ev) => Some(ev.columns_before_image()),
762            RowsEventData::PartialUpdateRowsEvent(ev) => Some(ev.columns_before_image()),
763        }
764    }
765
766    /// Returns columns in the after-image (only for WRITE and UPDATE).
767    ///
768    /// Each bit indicates whether corresponding column is used in the image.
769    pub fn columns_after_image(&'a self) -> Option<&'a BitSlice<u8>> {
770        match self {
771            RowsEventData::WriteRowsEventV1(ev) => Some(ev.columns_after_image()),
772            RowsEventData::UpdateRowsEventV1(ev) => Some(ev.columns_after_image()),
773            RowsEventData::DeleteRowsEventV1(_) => None,
774            RowsEventData::WriteRowsEvent(ev) => Some(ev.columns_after_image()),
775            RowsEventData::UpdateRowsEvent(ev) => Some(ev.columns_after_image()),
776            RowsEventData::DeleteRowsEvent(_) => None,
777            RowsEventData::PartialUpdateRowsEvent(ev) => Some(ev.columns_after_image()),
778        }
779    }
780
781    /// Returns raw rows data.
782    pub fn rows_data(&'a self) -> &'a [u8] {
783        match self {
784            RowsEventData::WriteRowsEventV1(ev) => ev.rows_data(),
785            RowsEventData::UpdateRowsEventV1(ev) => ev.rows_data(),
786            RowsEventData::DeleteRowsEventV1(ev) => ev.rows_data(),
787            RowsEventData::WriteRowsEvent(ev) => ev.rows_data(),
788            RowsEventData::UpdateRowsEvent(ev) => ev.rows_data(),
789            RowsEventData::DeleteRowsEvent(ev) => ev.rows_data(),
790            RowsEventData::PartialUpdateRowsEvent(ev) => ev.rows_data(),
791        }
792    }
793
794    /// Returns event flags.
795    pub fn flags(&self) -> RowsEventFlags {
796        match self {
797            RowsEventData::WriteRowsEventV1(ev) => ev.flags(),
798            RowsEventData::UpdateRowsEventV1(ev) => ev.flags(),
799            RowsEventData::DeleteRowsEventV1(ev) => ev.flags(),
800            RowsEventData::WriteRowsEvent(ev) => ev.flags(),
801            RowsEventData::UpdateRowsEvent(ev) => ev.flags(),
802            RowsEventData::DeleteRowsEvent(ev) => ev.flags(),
803            RowsEventData::PartialUpdateRowsEvent(ev) => ev.flags(),
804        }
805    }
806
807    /// Returns an iterator over event's rows given the corresponding `TableMapEvent`.
808    pub fn rows(&'a self, table_map_event: &'a TableMapEvent<'a>) -> RowsEventRows<'a> {
809        match self {
810            RowsEventData::WriteRowsEventV1(ev) => ev.rows(table_map_event),
811            RowsEventData::UpdateRowsEventV1(ev) => ev.rows(table_map_event),
812            RowsEventData::DeleteRowsEventV1(ev) => ev.rows(table_map_event),
813            RowsEventData::WriteRowsEvent(ev) => ev.rows(table_map_event),
814            RowsEventData::UpdateRowsEvent(ev) => ev.rows(table_map_event),
815            RowsEventData::DeleteRowsEvent(ev) => ev.rows(table_map_event),
816            RowsEventData::PartialUpdateRowsEvent(ev) => ev.rows(table_map_event),
817        }
818    }
819
820    pub fn into_owned(self) -> RowsEventData<'static> {
821        match self {
822            Self::WriteRowsEventV1(ev) => RowsEventData::WriteRowsEventV1(ev.into_owned()),
823            Self::UpdateRowsEventV1(ev) => RowsEventData::UpdateRowsEventV1(ev.into_owned()),
824            Self::DeleteRowsEventV1(ev) => RowsEventData::DeleteRowsEventV1(ev.into_owned()),
825            Self::WriteRowsEvent(ev) => RowsEventData::WriteRowsEvent(ev.into_owned()),
826            Self::UpdateRowsEvent(ev) => RowsEventData::UpdateRowsEvent(ev.into_owned()),
827            Self::DeleteRowsEvent(ev) => RowsEventData::DeleteRowsEvent(ev.into_owned()),
828            Self::PartialUpdateRowsEvent(ev) => {
829                RowsEventData::PartialUpdateRowsEvent(ev.into_owned())
830            }
831        }
832    }
833}
834
835impl MySerialize for RowsEventData<'_> {
836    fn serialize(&self, buf: &mut Vec<u8>) {
837        match self {
838            RowsEventData::WriteRowsEventV1(ev) => ev.serialize(buf),
839            RowsEventData::UpdateRowsEventV1(ev) => ev.serialize(buf),
840            RowsEventData::DeleteRowsEventV1(ev) => ev.serialize(buf),
841            RowsEventData::WriteRowsEvent(ev) => ev.serialize(buf),
842            RowsEventData::UpdateRowsEvent(ev) => ev.serialize(buf),
843            RowsEventData::DeleteRowsEvent(ev) => ev.serialize(buf),
844            RowsEventData::PartialUpdateRowsEvent(ev) => ev.serialize(buf),
845        }
846    }
847}