Skip to main content

mysql_common/binlog/
mod.rs

1// Copyright (c) 2020 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
9//! Binlog-related structures and functions. This implementation assumes
10//! binlog version >= 4 (MySql >= 5.0.0).
11//!
12//! All structures of this module contains raw data that may not necessarily be valid.
13//! Please consult the MySql documentation.
14
15// #![cfg(features = "binlog")]
16
17use std::{
18    collections::HashMap,
19    convert::TryFrom,
20    hash::Hash,
21    io::{
22        self, BufRead, Error,
23        ErrorKind::{InvalidData, UnexpectedEof},
24        Read, Write,
25    },
26};
27
28use crate::{
29    constants::ColumnType,
30    proto::{MyDeserialize, MySerialize},
31};
32
33#[allow(unused)]
34use self::events::TransactionPayloadEvent;
35
36use self::{
37    consts::{BinlogVersion, EventType},
38    events::{Event, FormatDescriptionEvent, RotateEvent, TableMapEvent},
39};
40
41pub mod consts;
42pub mod decimal;
43pub mod events;
44pub mod jsonb;
45pub mod jsondiff;
46pub mod misc;
47pub mod row;
48pub mod time;
49pub mod value;
50
51pub struct BinlogCtx<'a> {
52    pub event_size: usize,
53    pub fde: &'a FormatDescriptionEvent<'a>,
54    /// Raw event type byte from the binlog event header.
55    ///
56    /// This allows deserialize implementations to dispatch on the event type
57    /// when the same struct handles multiple wire formats (e.g. `GtidEvent`
58    /// handles both `GTID_EVENT` and `GTID_TAGGED_LOG_EVENT`).
59    pub event_type_raw: u8,
60}
61
62impl<'a> BinlogCtx<'a> {
63    pub fn new(event_size: usize, fde: &'a FormatDescriptionEvent<'a>, event_type_raw: u8) -> Self {
64        Self {
65            event_size,
66            fde,
67            event_type_raw,
68        }
69    }
70}
71
72/// Binlog event.
73pub trait BinlogStruct<'a>: MySerialize + MyDeserialize<'a, Ctx = BinlogCtx<'a>> {
74    /// Returns serialized length of this struct in bytes.
75    ///
76    /// *   implementation must truncate each field to its maximum length.
77    fn len(&self, version: BinlogVersion) -> usize;
78}
79
80pub trait BinlogEvent<'a>: BinlogStruct<'a> {
81    /// An event type, associated with this struct (if any).
82    const EVENT_TYPE: EventType;
83}
84
85/// A binlog file starts with a Binlog File Header `[ fe 'bin' ]`.
86#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
87pub struct BinlogFileHeader;
88
89impl BinlogFileHeader {
90    /// Length of a binlog file header.
91    pub const LEN: usize = 4;
92    /// Value of a binlog file header.
93    pub const VALUE: [u8; Self::LEN] = [0xfe, b'b', b'i', b'n'];
94
95    pub fn read<T: Read>(mut input: T) -> io::Result<Self> {
96        let mut buf = [0_u8; Self::LEN];
97        input.read_exact(&mut buf)?;
98
99        if buf != Self::VALUE {
100            return Err(Error::new(InvalidData, "invalid binlog file header"));
101        }
102
103        Ok(Self)
104    }
105
106    pub fn write<T: Write>(&self, _version: BinlogVersion, mut output: T) -> io::Result<()> {
107        output.write_all(&Self::VALUE)
108    }
109
110    pub fn len(&self, _version: BinlogVersion) -> usize {
111        Self::LEN
112    }
113}
114
115/// Reader for binlog events.
116///
117/// # Note
118///
119/// It's a low-level stream reader and only maintains actual fde and table map,
120/// so one should properly handle encountered events (see docs on [`EventStreamReader::read`]).
121#[derive(Debug)]
122pub struct EventStreamReader {
123    fde: FormatDescriptionEvent<'static>,
124    table_map: HashMap<u64, TableMapEvent<'static>>,
125}
126
127impl EventStreamReader {
128    /// Creates a new instance.
129    pub fn new(version: BinlogVersion) -> Self {
130        Self {
131            fde: FormatDescriptionEvent::new(version),
132            table_map: Default::default(),
133        }
134    }
135
136    /// Returns the format description event.
137    ///
138    /// Returns the default placeholder if there was no FDE yet.
139    pub fn get_fde(&self) -> &FormatDescriptionEvent<'static> {
140        &self.fde
141    }
142
143    /// Disable/Enable checksum verification without changing the original algorithm.
144    ///
145    /// See [`EventStreamReader::read_decompressed`].
146    pub(crate) fn set_checksum_enabled(&mut self, enabled: bool) {
147        self.fde.footer_mut().set_checksum_enabled(enabled);
148    }
149
150    /// Returns the table map event for the given table id.
151    ///
152    /// Should be available if rows event with this table id encountered in the stream.
153    pub fn get_tme(&self, table_id: u64) -> Option<&TableMapEvent<'static>> {
154        self.table_map.get(&table_id)
155    }
156
157    /// Will read next event from the given stream (Returns None if stream is exhausted).
158    ///
159    /// # Note
160    ///
161    /// Since MySql 8.0.20 it is possible for an event stream to contain an embedded
162    /// stream of events in form of a [`TransactionPayloadEvent`]. This means that
163    /// to properly handle table maps it is necessary to read the embedded stream
164    /// as soon as it is encountered (see [`EventStreamReader::read_decompressed`]).
165    pub fn read<T: BufRead>(&mut self, mut input: T) -> io::Result<Option<Event>> {
166        if input.fill_buf().map(|x| x.is_empty())? {
167            return Ok(None);
168        }
169
170        let event = Event::read(&self.fde, input)?;
171
172        self.handle_event(&event)?;
173
174        Ok(Some(event))
175    }
176
177    /// This function reads decompressed payload of a Transaction_payload_event
178    /// (see [`TransactionPayloadEvent::decompressed`]).
179    ///
180    /// The difference is that checksum verification will be disabled according to the WL#3549.
181    ///
182    /// # Warning
183    ///
184    /// This function can't be used to skip checksum verification for regular events.
185    ///
186    /// # Errors
187    ///
188    /// There is a list of events that should never be a part of a transaction payload and the list
189    /// includes the [`TransactionPayloadEvent`] itself.
190    /// This function will emit an [`io::ErrorKind::Other`] if [`TransactionPayloadEvent`]
191    /// is encountered within the compressed payload.
192    pub fn read_decompressed<T: BufRead>(&mut self, input: T) -> io::Result<Option<Event>> {
193        self.set_checksum_enabled(false);
194        let result = self.read(input);
195        self.set_checksum_enabled(true);
196        let Some(event) = result? else {
197            return Ok(None);
198        };
199
200        if event.header().event_type_raw() == EventType::TRANSACTION_PAYLOAD_EVENT as u8 {
201            return Err(io::Error::other("TRANSACTION_PAYLOAD_EVENT encountered"));
202        }
203
204        self.handle_event(&event)?;
205
206        Ok(Some(event))
207    }
208
209    fn handle_event(&mut self, event: &Event) -> io::Result<()> {
210        let event_type = event.header().event_type_raw();
211
212        if event_type == EventType::FORMAT_DESCRIPTION_EVENT as u8 {
213            // we'll redefine fde with an actual one
214            let fde = event.read_event::<FormatDescriptionEvent<'_>>()?;
215            self.fde = fde.into_owned().with_footer(event.footer());
216        } else if event_type == EventType::TABLE_MAP_EVENT as u8 {
217            // we'll maintain known table maps
218            let tme = event.read_event::<TableMapEvent<'_>>()?;
219            self.table_map.insert(tme.table_id(), tme.into_owned());
220        } else if event_type == EventType::ROTATE_EVENT as u8 {
221            // we'll keep table map size within reasonable bounds
222
223            // TODO: This value is arbitrary
224            const TABLE_MAP_MAX_SIZE: usize = 64;
225
226            let re = event.read_event::<RotateEvent<'_>>()?;
227            if !re.is_fake() {
228                self.table_map.clear();
229                self.table_map.shrink_to(TABLE_MAP_MAX_SIZE);
230            }
231        }
232
233        Ok(())
234    }
235}
236
237/// Binlog file.
238///
239/// It's an iterator over events in a binlog file.
240#[derive(Debug)]
241pub struct BinlogFile<T> {
242    reader: EventStreamReader,
243    read: T,
244}
245
246impl<T: BufRead> BinlogFile<T> {
247    /// Creates a new instance.
248    ///
249    /// It'll try to read binlog file header.
250    pub fn new(version: BinlogVersion, mut read: T) -> io::Result<Self> {
251        let reader = EventStreamReader::new(version);
252        BinlogFileHeader::read(&mut read)?;
253        Ok(Self { reader, read })
254    }
255
256    /// Returns a reference to the binlog stream reader.
257    pub fn reader(&self) -> &EventStreamReader {
258        &self.reader
259    }
260
261    /// Returns a mutable reference to the binlog stream reader.
262    pub fn reader_mut(&mut self) -> &mut EventStreamReader {
263        &mut self.reader
264    }
265}
266
267impl<T: BufRead> Iterator for BinlogFile<T> {
268    type Item = io::Result<Event>;
269
270    fn next(&mut self) -> Option<Self::Item> {
271        match self.reader.read(&mut self.read) {
272            Ok(event) => event.map(Ok),
273            Err(err) if err.kind() == UnexpectedEof => None,
274            Err(err) => Some(Err(err)),
275        }
276    }
277}
278
279impl ColumnType {
280    /// Returns type-specific metadata for this column type,
281    /// as well as the total number of occupied bytes.
282    ///
283    /// `is_array` must be true if `self` is from `MYSQL_TYPE_TYPED_ARRAY` metadata.
284    fn get_metadata<'a>(&self, ptr: &'a [u8], is_array: bool) -> Option<(&'a [u8], usize)> {
285        match self {
286            Self::MYSQL_TYPE_TINY_BLOB
287            | Self::MYSQL_TYPE_BLOB
288            | Self::MYSQL_TYPE_MEDIUM_BLOB
289            | Self::MYSQL_TYPE_LONG_BLOB
290            | Self::MYSQL_TYPE_DOUBLE
291            | Self::MYSQL_TYPE_FLOAT
292            | Self::MYSQL_TYPE_GEOMETRY
293            | Self::MYSQL_TYPE_TIME2
294            | Self::MYSQL_TYPE_DATETIME2
295            | Self::MYSQL_TYPE_TIMESTAMP2
296            | Self::MYSQL_TYPE_JSON
297            | Self::MYSQL_TYPE_VECTOR => ptr.get(..1).map(|x| (x, 1)),
298            Self::MYSQL_TYPE_VARCHAR => {
299                if is_array {
300                    ptr.get(..3).map(|x| (x, 3))
301                } else {
302                    ptr.get(..2).map(|x| (x, 2))
303                }
304            }
305            Self::MYSQL_TYPE_NEWDECIMAL
306            | Self::MYSQL_TYPE_SET
307            | Self::MYSQL_TYPE_ENUM
308            | Self::MYSQL_TYPE_STRING
309            | Self::MYSQL_TYPE_BIT => ptr.get(..2).map(|x| (x, 2)),
310            Self::MYSQL_TYPE_TYPED_ARRAY => Self::try_from(*ptr.first()?)
311                .ok()?
312                .get_metadata(ptr.get(1..)?, true)
313                .map(|(x, n)| (x, n + 1)),
314            Self::MYSQL_TYPE_DECIMAL
315            | Self::MYSQL_TYPE_TINY
316            | Self::MYSQL_TYPE_SHORT
317            | Self::MYSQL_TYPE_LONG
318            | Self::MYSQL_TYPE_NULL
319            | Self::MYSQL_TYPE_TIMESTAMP
320            | Self::MYSQL_TYPE_LONGLONG
321            | Self::MYSQL_TYPE_INT24
322            | Self::MYSQL_TYPE_DATE
323            | Self::MYSQL_TYPE_TIME
324            | Self::MYSQL_TYPE_DATETIME
325            | Self::MYSQL_TYPE_YEAR
326            | Self::MYSQL_TYPE_NEWDATE
327            | Self::MYSQL_TYPE_UNKNOWN
328            | Self::MYSQL_TYPE_VAR_STRING => Some((&[], 0)),
329        }
330    }
331}
332
333#[cfg(test)]
334mod tests {
335    use std::{
336        collections::HashMap,
337        convert::TryFrom,
338        io,
339        iter::{once, repeat},
340    };
341
342    use super::{
343        BinlogFile, BinlogFileHeader, BinlogVersion,
344        consts::{EventFlags, EventType},
345        events::{BinlogEventHeader, EventData, GtidEvent},
346    };
347
348    use crate::{
349        binlog::{
350            events::{OptionalMetadataField, RowsEventData},
351            value::BinlogValue,
352        },
353        collations::CollationId,
354        constants::ColumnFlags,
355        proto::MySerialize,
356        row::convert::from_row,
357        value::Value,
358    };
359
360    const BINLOG_FILE: &[u8] = &[
361        0xfe, 0x62, 0x69, 0x6e, 0xfc, 0x35, 0xbb, 0x4a, 0x0f, 0x01, 0x00, 0x00, 0x00, 0x5e, 0x00,
362        0x00, 0x00, 0x62, 0x00, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x35, 0x2e, 0x30, 0x2e, 0x38,
363        0x36, 0x2d, 0x64, 0x65, 0x62, 0x75, 0x67, 0x2d, 0x6c, 0x6f, 0x67, 0x00, 0x00, 0x00, 0x00,
364        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
365        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
366        0xfc, 0x35, 0xbb, 0x4a, 0x13, 0x38, 0x0d, 0x00, 0x08, 0x00, 0x12, 0x00, 0x04, 0x04, 0x04,
367        0x04, 0x12, 0x00, 0x00, 0x4b, 0x00, 0x04, 0x1a, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00,
368        0x00, 0x00, 0x64, 0x00, 0x00, 0x00, 0xc6, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00,
369        0x00, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00,
370        0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04,
371        0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x74, 0x65, 0x73, 0x74, 0x00, 0x63, 0x72, 0x65, 0x61,
372        0x74, 0x65, 0x20, 0x74, 0x61, 0x62, 0x6c, 0x65, 0x20, 0x74, 0x31, 0x28, 0x61, 0x20, 0x69,
373        0x6e, 0x74, 0x29, 0x20, 0x65, 0x6e, 0x67, 0x69, 0x6e, 0x65, 0x3d, 0x20, 0x69, 0x6e, 0x6e,
374        0x6f, 0x64, 0x62, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x65, 0x00, 0x00,
375        0x00, 0x2b, 0x01, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
376        0x05, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00,
377        0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08,
378        0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c, 0x00, 0x63, 0x72, 0x65, 0x61, 0x74, 0x65, 0x20, 0x74,
379        0x61, 0x62, 0x6c, 0x65, 0x20, 0x74, 0x32, 0x28, 0x61, 0x20, 0x69, 0x6e, 0x74, 0x29, 0x20,
380        0x65, 0x6e, 0x67, 0x69, 0x6e, 0x65, 0x3d, 0x20, 0x69, 0x6e, 0x6e, 0x6f, 0x64, 0x62, 0xfd,
381        0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x45, 0x00, 0x00, 0x00, 0x70, 0x01, 0x00,
382        0x00, 0x08, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00, 0x1a,
383        0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
384        0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79, 0x73,
385        0x71, 0x6c, 0x00, 0x42, 0x45, 0x47, 0x49, 0x4e, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00,
386        0x00, 0x00, 0x5c, 0x00, 0x00, 0x00, 0xcc, 0x01, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00,
387        0x00, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00,
388        0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04,
389        0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x74, 0x65, 0x73, 0x74, 0x00, 0x69, 0x6e, 0x73, 0x65,
390        0x72, 0x74, 0x20, 0x69, 0x6e, 0x74, 0x6f, 0x20, 0x74, 0x31, 0x20, 0x28, 0x61, 0x29, 0x20,
391        0x76, 0x61, 0x6c, 0x75, 0x65, 0x73, 0x20, 0x28, 0x31, 0x29, 0xfd, 0x35, 0xbb, 0x4a, 0x02,
392        0x01, 0x00, 0x00, 0x00, 0x5d, 0x00, 0x00, 0x00, 0x29, 0x02, 0x00, 0x00, 0x00, 0x00, 0x01,
393        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40,
394        0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74,
395        0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c, 0x00, 0x69,
396        0x6e, 0x73, 0x65, 0x72, 0x74, 0x20, 0x69, 0x6e, 0x74, 0x6f, 0x20, 0x74, 0x32, 0x20, 0x28,
397        0x61, 0x29, 0x20, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x73, 0x20, 0x28, 0x31, 0x29, 0xfd, 0x35,
398        0xbb, 0x4a, 0x10, 0x01, 0x00, 0x00, 0x00, 0x1b, 0x00, 0x00, 0x00, 0x44, 0x02, 0x00, 0x00,
399        0x00, 0x00, 0x0b, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0xfd, 0x35, 0xbb, 0x4a, 0x02,
400        0x01, 0x00, 0x00, 0x00, 0x64, 0x00, 0x00, 0x00, 0xa8, 0x02, 0x00, 0x00, 0x00, 0x00, 0x01,
401        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40,
402        0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74,
403        0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x74, 0x65, 0x73, 0x74, 0x00, 0x63, 0x72,
404        0x65, 0x61, 0x74, 0x65, 0x20, 0x74, 0x61, 0x62, 0x6c, 0x65, 0x20, 0x74, 0x33, 0x28, 0x61,
405        0x20, 0x69, 0x6e, 0x74, 0x29, 0x20, 0x65, 0x6e, 0x67, 0x69, 0x6e, 0x65, 0x3d, 0x20, 0x69,
406        0x6e, 0x6e, 0x6f, 0x64, 0x62, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x65,
407        0x00, 0x00, 0x00, 0x0d, 0x03, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00,
408        0x00, 0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00,
409        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08,
410        0x00, 0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c, 0x00, 0x63, 0x72, 0x65, 0x61, 0x74, 0x65,
411        0x20, 0x74, 0x61, 0x62, 0x6c, 0x65, 0x20, 0x74, 0x34, 0x28, 0x61, 0x20, 0x69, 0x6e, 0x74,
412        0x29, 0x20, 0x65, 0x6e, 0x67, 0x69, 0x6e, 0x65, 0x3d, 0x20, 0x6d, 0x79, 0x69, 0x73, 0x61,
413        0x6d, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x45, 0x00, 0x00, 0x00, 0x52,
414        0x03, 0x00, 0x00, 0x08, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00,
415        0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
416        0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d,
417        0x79, 0x73, 0x71, 0x6c, 0x00, 0x42, 0x45, 0x47, 0x49, 0x4e, 0xfd, 0x35, 0xbb, 0x4a, 0x02,
418        0x01, 0x00, 0x00, 0x00, 0x5c, 0x00, 0x00, 0x00, 0xae, 0x03, 0x00, 0x00, 0x00, 0x00, 0x01,
419        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40,
420        0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74,
421        0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x74, 0x65, 0x73, 0x74, 0x00, 0x69, 0x6e,
422        0x73, 0x65, 0x72, 0x74, 0x20, 0x69, 0x6e, 0x74, 0x6f, 0x20, 0x74, 0x33, 0x20, 0x28, 0x61,
423        0x29, 0x20, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x73, 0x20, 0x28, 0x32, 0x29, 0xfd, 0x35, 0xbb,
424        0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x5d, 0x00, 0x00, 0x00, 0x0b, 0x04, 0x00, 0x00, 0x00,
425        0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00,
426        0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03,
427        0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c,
428        0x00, 0x69, 0x6e, 0x73, 0x65, 0x72, 0x74, 0x20, 0x69, 0x6e, 0x74, 0x6f, 0x20, 0x74, 0x34,
429        0x20, 0x28, 0x61, 0x29, 0x20, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x73, 0x20, 0x28, 0x32, 0x29,
430        0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x48, 0x00, 0x00, 0x00, 0x53, 0x04,
431        0x00, 0x00, 0x08, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00,
432        0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
433        0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79,
434        0x73, 0x71, 0x6c, 0x00, 0x52, 0x4f, 0x4c, 0x4c, 0x42, 0x41, 0x43, 0x4b, 0xfd, 0x35, 0xbb,
435        0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x61, 0x00, 0x00, 0x00, 0xb4, 0x04, 0x00, 0x00, 0x00,
436        0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x1a, 0x00, 0x00,
437        0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03,
438        0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x74, 0x65, 0x73, 0x74, 0x00,
439        0x63, 0x72, 0x65, 0x61, 0x74, 0x65, 0x20, 0x74, 0x61, 0x62, 0x6c, 0x65, 0x20, 0x74, 0x35,
440        0x28, 0x61, 0x20, 0x69, 0x6e, 0x74, 0x29, 0x20, 0x65, 0x6e, 0x67, 0x69, 0x6e, 0x65, 0x3d,
441        0x20, 0x4e, 0x44, 0x42, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x62, 0x00,
442        0x00, 0x00, 0x16, 0x05, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
443        0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00,
444        0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00,
445        0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c, 0x00, 0x63, 0x72, 0x65, 0x61, 0x74, 0x65, 0x20,
446        0x74, 0x61, 0x62, 0x6c, 0x65, 0x20, 0x74, 0x36, 0x28, 0x61, 0x20, 0x69, 0x6e, 0x74, 0x29,
447        0x20, 0x65, 0x6e, 0x67, 0x69, 0x6e, 0x65, 0x3d, 0x20, 0x4e, 0x44, 0x42, 0xfd, 0x35, 0xbb,
448        0x4a, 0x02, 0x01, 0x00, 0x00, 0x00, 0x45, 0x00, 0x00, 0x00, 0x5b, 0x05, 0x00, 0x00, 0x08,
449        0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00,
450        0x00, 0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03,
451        0x73, 0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c,
452        0x00, 0x42, 0x45, 0x47, 0x49, 0x4e, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00, 0x00, 0x00,
453        0x5c, 0x00, 0x00, 0x00, 0xb7, 0x05, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00,
454        0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00, 0x01, 0x00,
455        0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04, 0x08, 0x00,
456        0x08, 0x00, 0x08, 0x00, 0x74, 0x65, 0x73, 0x74, 0x00, 0x69, 0x6e, 0x73, 0x65, 0x72, 0x74,
457        0x20, 0x69, 0x6e, 0x74, 0x6f, 0x20, 0x74, 0x35, 0x20, 0x28, 0x61, 0x29, 0x20, 0x76, 0x61,
458        0x6c, 0x75, 0x65, 0x73, 0x20, 0x28, 0x33, 0x29, 0xfd, 0x35, 0xbb, 0x4a, 0x02, 0x01, 0x00,
459        0x00, 0x00, 0x5d, 0x00, 0x00, 0x00, 0x14, 0x06, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00,
460        0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00, 0x40, 0x00, 0x00,
461        0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73, 0x74, 0x64, 0x04,
462        0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c, 0x00, 0x69, 0x6e, 0x73,
463        0x65, 0x72, 0x74, 0x20, 0x69, 0x6e, 0x74, 0x6f, 0x20, 0x74, 0x36, 0x20, 0x28, 0x61, 0x29,
464        0x20, 0x76, 0x61, 0x6c, 0x75, 0x65, 0x73, 0x20, 0x28, 0x33, 0x29, 0xfd, 0x35, 0xbb, 0x4a,
465        0x02, 0x01, 0x00, 0x00, 0x00, 0x46, 0x00, 0x00, 0x00, 0x5a, 0x06, 0x00, 0x00, 0x08, 0x00,
466        0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x05, 0x00, 0x00, 0x1a, 0x00, 0x00, 0x00,
467        0x40, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x00, 0x06, 0x03, 0x73,
468        0x74, 0x64, 0x04, 0x08, 0x00, 0x08, 0x00, 0x08, 0x00, 0x6d, 0x79, 0x73, 0x71, 0x6c, 0x00,
469        0x43, 0x4f, 0x4d, 0x4d, 0x49, 0x54, 0xfd, 0x35, 0xbb, 0x4a, 0x04, 0x01, 0x00, 0x00, 0x00,
470        0x2c, 0x00, 0x00, 0x00, 0x86, 0x06, 0x00, 0x00, 0x00, 0x00, 0x04, 0x00, 0x00, 0x00, 0x00,
471        0x00, 0x00, 0x00, 0x6d, 0x61, 0x73, 0x74, 0x65, 0x72, 0x2d, 0x62, 0x69, 0x6e, 0x2e, 0x30,
472        0x30, 0x30, 0x30, 0x30, 0x32,
473    ];
474
475    #[test]
476    fn binlog_file_header_roundtrip() -> io::Result<()> {
477        let mut output = Vec::new();
478
479        let binlog_file_header = BinlogFileHeader::read(BINLOG_FILE)?;
480        binlog_file_header.write(BinlogVersion::Version4, &mut output)?;
481
482        assert_eq!(&output[..], &BINLOG_FILE[..BinlogFileHeader::LEN]);
483
484        Ok(())
485    }
486
487    #[test]
488    fn binlog_file_iterator() -> io::Result<()> {
489        let binlog_file = BinlogFile::new(BinlogVersion::Version4, BINLOG_FILE)?;
490
491        let mut total = 0;
492        let mut ev_pos = 4;
493
494        for (i, ev) in binlog_file.enumerate() {
495            let data_start = ev_pos + BinlogEventHeader::LEN;
496            let ev = ev?;
497            match i {
498                0 => {
499                    assert_eq!(
500                        ev.header(),
501                        BinlogEventHeader::new(
502                            1253783036,
503                            EventType::FORMAT_DESCRIPTION_EVENT,
504                            1,
505                            94,
506                            98,
507                            EventFlags::empty()
508                        )
509                    )
510                }
511                1 => assert_eq!(
512                    ev.header(),
513                    BinlogEventHeader::new(
514                        1253783037,
515                        EventType::QUERY_EVENT,
516                        1,
517                        100,
518                        198,
519                        EventFlags::empty()
520                    )
521                ),
522                2 => assert_eq!(
523                    ev.header(),
524                    BinlogEventHeader::new(
525                        1253783037,
526                        EventType::QUERY_EVENT,
527                        1,
528                        101,
529                        299,
530                        EventFlags::empty()
531                    )
532                ),
533                3 => assert_eq!(
534                    ev.header(),
535                    BinlogEventHeader::new(
536                        1253783037,
537                        EventType::QUERY_EVENT,
538                        1,
539                        69,
540                        368,
541                        EventFlags::LOG_EVENT_SUPPRESS_USE_F
542                    )
543                ),
544                4 => assert_eq!(
545                    ev.header(),
546                    BinlogEventHeader::new(
547                        1253783037,
548                        EventType::QUERY_EVENT,
549                        1,
550                        92,
551                        460,
552                        EventFlags::empty()
553                    )
554                ),
555                5 => assert_eq!(
556                    ev.header(),
557                    BinlogEventHeader::new(
558                        1253783037,
559                        EventType::QUERY_EVENT,
560                        1,
561                        93,
562                        553,
563                        EventFlags::empty()
564                    )
565                ),
566                6 => assert_eq!(
567                    ev.header(),
568                    BinlogEventHeader::new(
569                        1253783037,
570                        EventType::XID_EVENT,
571                        1,
572                        27,
573                        580,
574                        EventFlags::empty()
575                    )
576                ),
577                7 => assert_eq!(
578                    ev.header(),
579                    BinlogEventHeader::new(
580                        1253783037,
581                        EventType::QUERY_EVENT,
582                        1,
583                        100,
584                        680,
585                        EventFlags::empty()
586                    )
587                ),
588                8 => assert_eq!(
589                    ev.header(),
590                    BinlogEventHeader::new(
591                        1253783037,
592                        EventType::QUERY_EVENT,
593                        1,
594                        101,
595                        781,
596                        EventFlags::empty()
597                    )
598                ),
599                9 => assert_eq!(
600                    ev.header(),
601                    BinlogEventHeader::new(
602                        1253783037,
603                        EventType::QUERY_EVENT,
604                        1,
605                        69,
606                        850,
607                        EventFlags::LOG_EVENT_SUPPRESS_USE_F
608                    )
609                ),
610                10 => assert_eq!(
611                    ev.header(),
612                    BinlogEventHeader::new(
613                        1253783037,
614                        EventType::QUERY_EVENT,
615                        1,
616                        92,
617                        942,
618                        EventFlags::empty()
619                    )
620                ),
621                11 => assert_eq!(
622                    ev.header(),
623                    BinlogEventHeader::new(
624                        1253783037,
625                        EventType::QUERY_EVENT,
626                        1,
627                        93,
628                        1035,
629                        EventFlags::empty()
630                    )
631                ),
632                12 => assert_eq!(
633                    ev.header(),
634                    BinlogEventHeader::new(
635                        1253783037,
636                        EventType::QUERY_EVENT,
637                        1,
638                        72,
639                        1107,
640                        EventFlags::LOG_EVENT_SUPPRESS_USE_F
641                    )
642                ),
643                13 => assert_eq!(
644                    ev.header(),
645                    BinlogEventHeader::new(
646                        1253783037,
647                        EventType::QUERY_EVENT,
648                        1,
649                        97,
650                        1204,
651                        EventFlags::empty()
652                    )
653                ),
654                14 => assert_eq!(
655                    ev.header(),
656                    BinlogEventHeader::new(
657                        1253783037,
658                        EventType::QUERY_EVENT,
659                        1,
660                        98,
661                        1302,
662                        EventFlags::empty()
663                    )
664                ),
665                15 => assert_eq!(
666                    ev.header(),
667                    BinlogEventHeader::new(
668                        1253783037,
669                        EventType::QUERY_EVENT,
670                        1,
671                        69,
672                        1371,
673                        EventFlags::LOG_EVENT_SUPPRESS_USE_F
674                    )
675                ),
676                16 => assert_eq!(
677                    ev.header(),
678                    BinlogEventHeader::new(
679                        1253783037,
680                        EventType::QUERY_EVENT,
681                        1,
682                        92,
683                        1463,
684                        EventFlags::empty()
685                    )
686                ),
687                17 => assert_eq!(
688                    ev.header(),
689                    BinlogEventHeader::new(
690                        1253783037,
691                        EventType::QUERY_EVENT,
692                        1,
693                        93,
694                        1556,
695                        EventFlags::empty()
696                    )
697                ),
698                18 => assert_eq!(
699                    ev.header(),
700                    BinlogEventHeader::new(
701                        1253783037,
702                        EventType::QUERY_EVENT,
703                        1,
704                        70,
705                        1626,
706                        EventFlags::LOG_EVENT_SUPPRESS_USE_F
707                    )
708                ),
709                19 => assert_eq!(
710                    ev.header(),
711                    BinlogEventHeader::new(
712                        1253783037,
713                        EventType::ROTATE_EVENT,
714                        1,
715                        44,
716                        1670,
717                        EventFlags::empty()
718                    )
719                ),
720                _ => panic!("too many"),
721            }
722
723            assert_eq!(
724                ev.data(),
725                &BINLOG_FILE[data_start
726                    ..(data_start + ev.header().event_size() as usize - BinlogEventHeader::LEN)],
727            );
728
729            total += 1;
730            ev_pos = ev.header().log_pos() as usize;
731        }
732
733        assert_eq!(total, 20);
734        Ok(())
735    }
736
737    #[test]
738    fn binlog_event_roundtrip() -> io::Result<()> {
739        const PATH: &str = "./test-data/binlogs";
740
741        let binlogs = std::fs::read_dir(PATH)?
742            .filter_map(|path| path.ok())
743            .map(|entry| entry.path())
744            .filter(|path| path.file_name().is_some());
745
746        'outer: for file_path in binlogs {
747            let file_data = std::fs::read(dbg!(&file_path))?;
748            let mut binlog_file = BinlogFile::new(BinlogVersion::Version4, &file_data[..])?;
749
750            let mut i = 0;
751            let mut ev_pos = 4;
752            let mut table_map_events = HashMap::new();
753
754            while let Some(ev) = binlog_file.next() {
755                i += 1;
756                let ev = ev?;
757                let _ = dbg!(ev.header().event_type());
758                let ev_end = ev_pos + ev.header().event_size() as usize;
759                let binlog_version = binlog_file.reader.fde.binlog_version();
760
761                let mut output = Vec::new();
762                ev.write(binlog_version, &mut output)?;
763
764                let event = match ev.read_data() {
765                    Ok(event) => {
766                        let event = match event {
767                            Some(e) => e,
768                            None => {
769                                if file_path.file_name().unwrap() == "mariadb-bin.000001" {
770                                    continue;
771                                } else {
772                                    dbg!(&ev);
773                                    panic!();
774                                }
775                            }
776                        };
777                        match event {
778                            EventData::TableMapEvent(ref ev) => {
779                                // store table maps for later use
780                                table_map_events.insert(ev.table_id(), ev.clone().into_owned());
781
782                                event
783                            }
784                            EventData::RowsEvent(ref rows_event) => {
785                                // iterate rows in a rows event
786                                let table_map_event =
787                                    binlog_file.reader().get_tme(rows_event.table_id()).unwrap();
788                                for row in rows_event.rows(table_map_event) {
789                                    let _row = row.unwrap();
790                                    if file_path.file_name().unwrap() == "mariadb-bin.000001" {
791                                        // should parse metadata for `binlog_row_metadata=FULL`
792                                        let after = _row.1.as_ref().unwrap();
793                                        let columns = after.columns_ref();
794
795                                        for col in columns.iter() {
796                                            assert_eq!(col.schema_ref(), b"toddy_test");
797                                            assert_eq!(col.table_ref(), b"outbox");
798                                            assert_eq!(col.org_table_ref(), b"outbox");
799                                        }
800
801                                        for (col, col_name) in columns.iter().zip([
802                                            "id",
803                                            "topic",
804                                            "event_type",
805                                            "event",
806                                            "created",
807                                        ]) {
808                                            assert_eq!(col.name_ref(), col_name.as_bytes());
809                                        }
810
811                                        for (col, f) in columns.iter().zip(
812                                            once(ColumnFlags::PRI_KEY_FLAG)
813                                                .chain(repeat(ColumnFlags::empty())),
814                                        ) {
815                                            assert_eq!(col.flags(), f);
816                                        }
817
818                                        for (col, charset) in columns.iter().zip([
819                                            CollationId::UNKNOWN_COLLATION_ID,
820                                            CollationId::UTF8MB4_GENERAL_CI,
821                                            CollationId::UTF8MB4_GENERAL_CI,
822                                            CollationId::BINARY,
823                                            CollationId::UNKNOWN_COLLATION_ID,
824                                        ]) {
825                                            assert_eq!(col.character_set(), charset as u16);
826                                        }
827                                    }
828                                }
829
830                                event
831                            }
832                            _ => event,
833                        }
834                    }
835                    Err(err)
836                        if err.kind() == std::io::ErrorKind::Other
837                            && ev.header().event_type() == Ok(EventType::XID_EVENT)
838                            && ev.header().event_size() == 0x26
839                            && file_path.file_name().unwrap() == "ver_5_1-wl2325_r.001" =>
840                    {
841                        // ver_5_1-wl2325_r.001 testfile contains broken xid event.
842                        continue 'outer;
843                    }
844                    Err(err)
845                        if err.kind() == std::io::ErrorKind::UnexpectedEof
846                            && ev.header().event_type() == Ok(EventType::QUERY_EVENT)
847                            && ev.header().event_size() == 171
848                            && file_path.file_name().unwrap() == "corrupt-relay-bin.000624" =>
849                    {
850                        // corrupt-relay-bin.000624 testfile contains broken query event.
851                        continue 'outer;
852                    }
853                    other => other.transpose().unwrap()?,
854                };
855
856                if file_path.file_name().unwrap() == "binlog-invisible-columns.000001"
857                    && let Some(EventData::TableMapEvent(ev)) = ev.read_data().unwrap()
858                {
859                    let optional_meta = ev.iter_optional_meta();
860                    for meta in optional_meta {
861                        meta.unwrap();
862                    }
863                }
864
865                if file_path.file_name().unwrap() == "json-opaque.binlog" {
866                    let event_data = ev.read_data().unwrap();
867
868                    /// Extracts first column of the binlog row after-image as a Jsonb::Value
869                    /// then parses it into the structured representation and compares with
870                    /// the expected value.
871                    macro_rules! extract_cmp {
872                        ($row:expr_2021, $expected:tt) => {
873                            let mut after = $row.1.unwrap().unwrap();
874                            let a = dbg!(after.pop().unwrap());
875                            let super::value::BinlogValue::Jsonb(a) = a else {
876                                panic!("BinlogValue::Jsonb(_) expected");
877                            };
878                            assert_eq!(
879                                serde_json::json!($expected),
880                                serde_json::Value::from(a.parse().unwrap())
881                            );
882                        };
883                    }
884
885                    match event_data {
886                        Some(EventData::RowsEvent(ev)) if i == 10 => {
887                            let table_map_event =
888                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
889                            let mut rows = ev.rows(table_map_event);
890                            extract_cmp!(rows.next().unwrap().unwrap(), {"a": "base64:type15:VQ=="});
891                        }
892                        Some(EventData::RowsEvent(ev)) if i == 12 => {
893                            let table_map_event =
894                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
895                            let mut rows = ev.rows(table_map_event);
896                            extract_cmp!(rows.next().unwrap().unwrap(), {"b": "2012-03-18"});
897                        }
898                        Some(EventData::RowsEvent(ev)) if i == 14 => {
899                            let table_map_event =
900                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
901                            let mut rows = ev.rows(table_map_event);
902                            extract_cmp!(rows.next().unwrap().unwrap(), {"c": "2012-03-18 11:30:45.000000"});
903                        }
904                        Some(EventData::RowsEvent(ev)) if i == 16 => {
905                            let table_map_event =
906                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
907                            let mut rows = ev.rows(table_map_event);
908                            extract_cmp!(rows.next().unwrap().unwrap(), {"c": "87:31:46.654321"});
909                        }
910                        Some(EventData::RowsEvent(ev)) if i == 18 => {
911                            let table_map_event =
912                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
913                            let mut rows = ev.rows(table_map_event);
914                            extract_cmp!(rows.next().unwrap().unwrap(), {"d": "123.456"});
915                        }
916                        Some(EventData::RowsEvent(ev)) if i == 20 => {
917                            let table_map_event =
918                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
919                            let mut rows = ev.rows(table_map_event);
920                            extract_cmp!(rows.next().unwrap().unwrap(), {"e": "9.00"});
921                        }
922                        Some(EventData::RowsEvent(ev)) if i == 22 => {
923                            let table_map_event =
924                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
925                            let mut rows = ev.rows(table_map_event);
926                            extract_cmp!(rows.next().unwrap().unwrap(), {"e": [0, 1, true, false]});
927                        }
928                        Some(EventData::RowsEvent(ev)) if i == 24 => {
929                            let table_map_event =
930                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
931                            let mut rows = ev.rows(table_map_event);
932                            extract_cmp!(rows.next().unwrap().unwrap(), {"e": null});
933                        }
934                        Some(EventData::RowsEvent(ev)) => {
935                            panic!("no more events expected i={}, {:?}", i, ev);
936                        }
937                        _ => (),
938                    }
939                }
940
941                if file_path.file_name().unwrap() == "vector.binlog" {
942                    let event_data = ev.read_data().unwrap();
943                    match event_data {
944                        Some(EventData::TableMapEvent(ev)) => {
945                            let optional_meta = ev.iter_optional_meta();
946                            match ev.table_name().as_ref() {
947                                "foo" => {
948                                    for meta in optional_meta {
949                                        if let OptionalMetadataField::Dimensionality(x) =
950                                            meta.unwrap()
951                                        {
952                                            assert_eq!(
953                                                x.iter_dimensionalities()
954                                                    .collect::<Result<Vec<_>, _>>()
955                                                    .unwrap(),
956                                                vec![3],
957                                            )
958                                        }
959                                    }
960                                }
961                                "bar" => {
962                                    for meta in optional_meta {
963                                        if let OptionalMetadataField::Dimensionality(x) =
964                                            meta.unwrap()
965                                        {
966                                            assert_eq!(
967                                                x.iter_dimensionalities()
968                                                    .collect::<Result<Vec<_>, _>>()
969                                                    .unwrap(),
970                                                vec![2, 4],
971                                            )
972                                        }
973                                    }
974                                }
975                                _ => (),
976                            }
977                        }
978                        Some(EventData::RowsEvent(ev)) if i == 12 => {
979                            let table_map_event =
980                                binlog_file.reader().get_tme(ev.table_id()).unwrap();
981                            let mut rows = ev.rows(table_map_event);
982
983                            let (None, Some(after)) = rows.next().unwrap().unwrap() else {
984                                panic!("Unexpected data");
985                            };
986                            let (id, vector_column): (u8, Vec<u8>) =
987                                from_row(crate::Row::try_from(after).unwrap());
988                            assert_eq!(id, 1);
989                            assert_eq!(
990                                vector_column,
991                                vec![205, 204, 140, 63, 205, 204, 12, 64, 51, 51, 83, 64]
992                            );
993
994                            let (None, Some(after)) = rows.next().unwrap().unwrap() else {
995                                panic!("Unexpected data");
996                            };
997                            let (id, vector_column): (u8, Vec<u8>) =
998                                from_row(crate::Row::try_from(after).unwrap());
999                            assert_eq!(id, 2);
1000                            assert_eq!(
1001                                vector_column,
1002                                vec![0, 0, 128, 63, 0, 0, 128, 191, 0, 0, 0, 0]
1003                            );
1004                        }
1005                        _ => (),
1006                    }
1007                }
1008
1009                if file_path.file_name().unwrap() == "mysql-enum-string-set.000001"
1010                    && let Some(EventData::RowsEvent(data)) = ev.read_data().unwrap()
1011                {
1012                    let table_map_event = binlog_file.reader().get_tme(data.table_id()).unwrap();
1013                    for row in data.rows(table_map_event) {
1014                        let (before, after) = row.unwrap();
1015                        match data {
1016                            RowsEventData::WriteRowsEvent(_) => {
1017                                assert!(before.is_none());
1018                                let after = after.unwrap().unwrap();
1019                                let mut j = 0;
1020                                for v in after {
1021                                    j += 1;
1022                                    match j {
1023                                            1 => assert_eq!(v, BinlogValue::Value("0123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789".into())),
1024                                            2 => assert_eq!(v, BinlogValue::Value("0123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789".into())),
1025                                            3 => assert_eq!(v, BinlogValue::Value(1_i8.into())),
1026                                            4 => assert_eq!(v, BinlogValue::Value([0b00000101_u8].into())),
1027                                            5 => assert_eq!(v, BinlogValue::Value("0123456789".into())),
1028
1029                                            _ => panic!(),
1030                                        }
1031                                }
1032                                assert_eq!(j, 5);
1033                            }
1034                            RowsEventData::UpdateRowsEvent(_) => {
1035                                let before = before.unwrap().unwrap();
1036                                let mut j = 0;
1037                                for v in before {
1038                                    j += 1;
1039                                    match j {
1040                                            1 => assert_eq!(v, BinlogValue::Value("0123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789".into())),
1041                                            2 => assert_eq!(v, BinlogValue::Value("0123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789".into())),
1042                                            3 => assert_eq!(v, BinlogValue::Value(1_i8.into())),
1043                                            4 => assert_eq!(v, BinlogValue::Value([0b00000101_u8].into())),
1044                                            5 => assert_eq!(v, BinlogValue::Value("0123456789".into())),
1045
1046                                            _ => panic!(),
1047                                        }
1048                                }
1049                                assert_eq!(j, 5);
1050
1051                                let after = after.unwrap().unwrap();
1052                                let mut j = 0;
1053                                for v in after {
1054                                    j += 1;
1055                                    match j {
1056                                            1 => assert_eq!(v, BinlogValue::Value("field1".into())),
1057                                            2 => assert_eq!(v, BinlogValue::Value("field_2".into())),
1058                                            3 => assert_eq!(v, BinlogValue::Value(2_i8.into())),
1059                                            4 => assert_eq!(v, BinlogValue::Value([0b00001010_u8].into())),
1060                                            5 => assert_eq!(v, BinlogValue::Value("0123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789".into())),
1061                                            _ => panic!(),
1062                                        }
1063                                }
1064                                assert_eq!(j, 5);
1065                            }
1066                            RowsEventData::DeleteRowsEvent(_) => {
1067                                assert!(after.is_none());
1068
1069                                let before = before.unwrap().unwrap();
1070                                let mut j = 0;
1071                                for v in before {
1072                                    j += 1;
1073                                    match j {
1074                                            1 => assert_eq!(v, BinlogValue::Value("field1".into())),
1075                                            2 => assert_eq!(v, BinlogValue::Value("field_2".into())),
1076                                            3 => assert_eq!(v, BinlogValue::Value(2_i8.into())),
1077                                            4 => assert_eq!(v, BinlogValue::Value([0b00001010_u8].into())),
1078                                            5 => assert_eq!(v, BinlogValue::Value("0123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456789012345678901234567890123456780123456789012345678901234567890123456789".into())),
1079                                            _ => panic!(),
1080                                        }
1081                                }
1082                                assert_eq!(j, 5);
1083                            }
1084                            _ => panic!(),
1085                        }
1086                    }
1087                }
1088
1089                if file_path.file_name().unwrap() == "mariadb-bin.000001" {
1090                    // Extraneous bytes in RotateEvent file name
1091                    // https://github.com/blackbeam/mysql_async/issues/189
1092                    if let Some(EventData::RotateEvent(ev)) = ev.read_data().unwrap() {
1093                        assert_ne!(ev.name_raw(), b"mariadb-bin.000001");
1094                    }
1095                }
1096
1097                if file_path.file_name().unwrap() == "mysql_type_bit.000001"
1098                    && let Some(EventData::RowsEvent(ev)) = ev.read_data().unwrap()
1099                {
1100                    let table_map_event = binlog_file.reader().get_tme(ev.table_id()).unwrap();
1101                    for row in ev.rows(table_map_event) {
1102                        let (before, after) = row.unwrap();
1103                        assert_eq!(before, None);
1104                        assert_eq!(
1105                            after.unwrap().unwrap(),
1106                            vec![
1107                                BinlogValue::Value(Value::Bytes(vec![0b100])),
1108                                BinlogValue::Value(Value::Bytes(b"foo".to_vec())),
1109                                BinlogValue::Value(Value::Bytes(vec![0b100000])),
1110                            ],
1111                        );
1112                    }
1113                }
1114
1115                if file_path.file_name().unwrap() != "mariadb-bin.000001" {
1116                    assert_eq!(output, &file_data[ev_pos..ev_end]);
1117                }
1118
1119                if file_path.file_name().unwrap() == "transaction_compression.000001"
1120                    && let Some(EventData::TransactionPayloadEvent(data)) = ev.read_data().unwrap()
1121                {
1122                    let mut payload = data.decompressed()?;
1123                    let reader = binlog_file.reader_mut();
1124
1125                    let mut binlog_ev = reader.read_decompressed(&mut payload)?.unwrap();
1126                    assert_eq!(binlog_ev.header().event_type(), Ok(EventType::QUERY_EVENT));
1127
1128                    binlog_ev = reader.read_decompressed(&mut payload)?.unwrap();
1129                    assert_eq!(
1130                        binlog_ev.header().event_type(),
1131                        Ok(EventType::TABLE_MAP_EVENT)
1132                    );
1133
1134                    binlog_ev = reader.read_decompressed(&mut payload)?.unwrap();
1135                    assert_eq!(
1136                        binlog_ev.header().event_type(),
1137                        Ok(EventType::WRITE_ROWS_EVENT)
1138                    );
1139
1140                    binlog_ev = reader.read_decompressed(&mut payload)?.unwrap();
1141                    assert_eq!(binlog_ev.header().event_type(), Ok(EventType::XID_EVENT));
1142                    assert!(reader.read_decompressed(&mut payload)?.is_none());
1143                }
1144                output = Vec::new();
1145                event.serialize(&mut output);
1146
1147                if matches!(event, EventData::UserVarEvent(_)) {
1148                    // Server may or may not write the flags field, but we will always write it.
1149                    assert_eq!(&output[..ev.data().len()], ev.data());
1150                    assert!(output.len() == ev.data().len() || output.len() == ev.data().len() + 1);
1151                } else if (matches!(event, EventData::GtidEvent(_))
1152                    || matches!(event, EventData::AnonymousGtidEvent(_)))
1153                    && ev.fde().split_version() < (5, 7, 0)
1154                {
1155                    // MySql 5.6 does not write TS_TYPE and following post-header fields
1156                    assert_eq!(&output[..GtidEvent::POST_HEADER_LENGTH - 1 - 16], ev.data());
1157                } else if (matches!(event, EventData::GtidEvent(_))
1158                    || matches!(event, EventData::AnonymousGtidEvent(_)))
1159                    && ev.fde().split_version() < (5, 8, 0)
1160                {
1161                    // MySql 5.7 contains only post-header in this event
1162                    assert_eq!(&output[..GtidEvent::POST_HEADER_LENGTH], ev.data());
1163                } else {
1164                    assert_eq!(output, ev.data());
1165                }
1166
1167                // https://github.com/blackbeam/rust_mysql_common/issues/162
1168                if file_path.file_name().unwrap() == "minimal_row_metadata.000001" {
1169                    let event_data = ev.read_data().unwrap();
1170                    if let Some(EventData::RowsEvent(ev)) = event_data {
1171                        let table_map_event = binlog_file.reader().get_tme(ev.table_id()).unwrap();
1172                        let row = ev.rows(table_map_event).next().unwrap().unwrap();
1173                        let after_image = row.1.unwrap();
1174                        after_image
1175                            .columns()
1176                            .iter()
1177                            .enumerate()
1178                            .for_each(|(i, column)| {
1179                                if column.name_str() == "@2" {
1180                                    assert_eq!(
1181                                        column.character_set(),
1182                                        CollationId::UTF8MB4_0900_AI_CI as u16
1183                                    );
1184                                } else if column.name_str() == "@4" {
1185                                    match after_image.as_ref(i).unwrap() {
1186                                        BinlogValue::Value(val) => {
1187                                            assert_eq!(val, &Value::Int(3230202323));
1188                                        }
1189                                        _ => panic!("Expected a value"),
1190                                    }
1191                                }
1192                            });
1193                    }
1194                }
1195
1196                if file_path.file_name().unwrap() == "time_issue.000001" {
1197                    let event_data = ev.read_data().unwrap();
1198                    if let Some(EventData::RowsEvent(ev)) = event_data {
1199                        let table_map_event = binlog_file.reader().get_tme(ev.table_id()).unwrap();
1200                        let row = ev.rows(table_map_event).next().unwrap().unwrap();
1201                        let after_image = row.1.unwrap();
1202                        after_image.columns().iter().enumerate().for_each(
1203                            |(i, _)| match after_image.as_ref(i).unwrap() {
1204                                BinlogValue::Value(val) => {
1205                                    assert_eq!(val, &Value::Time(true, 21, 3, 48, 27, 0));
1206                                }
1207                                _ => panic!("Expected a value"),
1208                            },
1209                        );
1210                    }
1211                }
1212                ev_pos = ev_end;
1213            }
1214        }
1215
1216        Ok(())
1217    }
1218
1219    #[test]
1220    fn gtid_tagged_log_event_from_binlog() -> io::Result<()> {
1221        use crate::packets::GnoInterval;
1222
1223        let file_data =
1224            std::fs::read("./test-data/binlogs/binlog_transaction_with_GTID_TAG.000001")?;
1225        let binlog_file = BinlogFile::new(BinlogVersion::Version4, &file_data[..])?;
1226
1227        let mut found_tagged_gtid = false;
1228        let mut found_previous_gtids = false;
1229
1230        // UUID: 55778904-0299-11f1-b1b8-4ef0c4956feb
1231        let expected_sid = [
1232            0x55, 0x77, 0x89, 0x04, // 55778904
1233            0x02, 0x99, // 0299
1234            0x11, 0xf1, // 11f1
1235            0xb1, 0xb8, // b1b8
1236            0x4e, 0xf0, 0xc4, 0x95, 0x6f, 0xeb, // 4ef0c4956feb
1237        ];
1238
1239        for ev in binlog_file {
1240            let ev = ev?;
1241            if ev.header().event_type() == Ok(EventType::PREVIOUS_GTIDS_EVENT) {
1242                let event_data = ev.read_data()?.expect("should parse event data");
1243                match event_data {
1244                    EventData::PreviousGtidsEvent(prev_ev) => {
1245                        let sids = prev_ev.sids();
1246                        assert_eq!(sids.len(), 2, "expected 2 SID entries");
1247
1248                        // Entry 1: no tag, interval [1, 14)
1249                        assert_eq!(sids[0].uuid(), expected_sid);
1250                        assert!(sids[0].tag().is_none(), "first entry should have no tag");
1251                        assert_eq!(sids[0].intervals(), &[GnoInterval::new(1, 14)]);
1252
1253                        // Entry 2: tag "mytag", interval [1, 3)
1254                        assert_eq!(sids[1].uuid(), expected_sid);
1255                        assert_eq!(sids[1].tag().map(|t| t.as_str()), Some("mytag"));
1256                        assert_eq!(sids[1].intervals(), &[GnoInterval::new(1, 3)]);
1257
1258                        found_previous_gtids = true;
1259                    }
1260                    other => panic!(
1261                        "expected PreviousGtidsEvent, got {:?}",
1262                        std::mem::discriminant(&other)
1263                    ),
1264                }
1265            }
1266            if ev.header().event_type() == Ok(EventType::GTID_TAGGED_LOG_EVENT) {
1267                let event_data = ev.read_data()?.expect("should parse event data");
1268                match event_data {
1269                    EventData::GtidEvent(gtid) => {
1270                        assert!(gtid.is_tagged());
1271                        assert_eq!(gtid.tag().unwrap().as_str(), "mytag");
1272                        assert_eq!(gtid.sid(), expected_sid);
1273                        assert_eq!(gtid.gno(), 3);
1274                        found_tagged_gtid = true;
1275                    }
1276                    other => panic!(
1277                        "expected GtidEvent (tagged), got {:?}",
1278                        std::mem::discriminant(&other)
1279                    ),
1280                }
1281            }
1282        }
1283
1284        assert!(
1285            found_tagged_gtid,
1286            "GTID_TAGGED_LOG_EVENT not found in binlog"
1287        );
1288        assert!(
1289            found_previous_gtids,
1290            "PREVIOUS_GTIDS_EVENT not found in binlog"
1291        );
1292        Ok(())
1293    }
1294
1295    #[test]
1296    fn untagged_previous_gtids_event_from_binlog() -> io::Result<()> {
1297        use crate::packets::GnoInterval;
1298
1299        let file_data =
1300            std::fs::read("./test-data/binlogs/binlog_transaction_previous_GTID_no_tag.000001")?;
1301        let binlog_file = BinlogFile::new(BinlogVersion::Version4, &file_data[..])?;
1302
1303        let mut found = false;
1304
1305        // UUID: b9b88c66-0755-11f1-9899-4a9da94c4d71
1306        let expected_uuid = [
1307            0xb9, 0xb8, 0x8c, 0x66, // b9b88c66
1308            0x07, 0x55, // 0755
1309            0x11, 0xf1, // 11f1
1310            0x98, 0x99, // 9899
1311            0x4a, 0x9d, 0xa9, 0x4c, 0x4d, 0x71, // 4a9da94c4d71
1312        ];
1313
1314        for ev in binlog_file {
1315            let ev = ev?;
1316            if ev.header().event_type() == Ok(EventType::PREVIOUS_GTIDS_EVENT) {
1317                let event_data = ev.read_data()?.expect("should parse event data");
1318                match event_data {
1319                    EventData::PreviousGtidsEvent(prev_ev) => {
1320                        let sids = prev_ev.sids();
1321                        assert_eq!(sids.len(), 1);
1322
1323                        assert_eq!(sids[0].uuid(), expected_uuid);
1324                        assert!(sids[0].tag().is_none());
1325                        assert_eq!(sids[0].intervals(), &[GnoInterval::new(1, 3)]);
1326
1327                        found = true;
1328                    }
1329                    other => panic!(
1330                        "expected PreviousGtidsEvent, got {:?}",
1331                        std::mem::discriminant(&other)
1332                    ),
1333                }
1334            }
1335        }
1336
1337        assert!(found, "PREVIOUS_GTIDS_EVENT not found in binlog");
1338        Ok(())
1339    }
1340}