1use 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 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
72pub trait BinlogStruct<'a>: MySerialize + MyDeserialize<'a, Ctx = BinlogCtx<'a>> {
74 fn len(&self, version: BinlogVersion) -> usize;
78}
79
80pub trait BinlogEvent<'a>: BinlogStruct<'a> {
81 const EVENT_TYPE: EventType;
83}
84
85#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
87pub struct BinlogFileHeader;
88
89impl BinlogFileHeader {
90 pub const LEN: usize = 4;
92 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#[derive(Debug)]
122pub struct EventStreamReader {
123 fde: FormatDescriptionEvent<'static>,
124 table_map: HashMap<u64, TableMapEvent<'static>>,
125}
126
127impl EventStreamReader {
128 pub fn new(version: BinlogVersion) -> Self {
130 Self {
131 fde: FormatDescriptionEvent::new(version),
132 table_map: Default::default(),
133 }
134 }
135
136 pub fn get_fde(&self) -> &FormatDescriptionEvent<'static> {
140 &self.fde
141 }
142
143 pub(crate) fn set_checksum_enabled(&mut self, enabled: bool) {
147 self.fde.footer_mut().set_checksum_enabled(enabled);
148 }
149
150 pub fn get_tme(&self, table_id: u64) -> Option<&TableMapEvent<'static>> {
154 self.table_map.get(&table_id)
155 }
156
157 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 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 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 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 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#[derive(Debug)]
241pub struct BinlogFile<T> {
242 reader: EventStreamReader,
243 read: T,
244}
245
246impl<T: BufRead> BinlogFile<T> {
247 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 pub fn reader(&self) -> &EventStreamReader {
258 &self.reader
259 }
260
261 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 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 table_map_events.insert(ev.table_id(), ev.clone().into_owned());
781
782 event
783 }
784 EventData::RowsEvent(ref rows_event) => {
785 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 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 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 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 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 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 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 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 assert_eq!(&output[..GtidEvent::POST_HEADER_LENGTH], ev.data());
1163 } else {
1164 assert_eq!(output, ev.data());
1165 }
1166
1167 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 let expected_sid = [
1232 0x55, 0x77, 0x89, 0x04, 0x02, 0x99, 0x11, 0xf1, 0xb1, 0xb8, 0x4e, 0xf0, 0xc4, 0x95, 0x6f, 0xeb, ];
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 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 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 let expected_uuid = [
1307 0xb9, 0xb8, 0x8c, 0x66, 0x07, 0x55, 0x11, 0xf1, 0x98, 0x99, 0x4a, 0x9d, 0xa9, 0x4c, 0x4d, 0x71, ];
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}