1use byteorder::{ByteOrder, LittleEndian};
7use serde::Serialize;
8
9use super::event::BinlogEventType;
10
11#[derive(Debug, Clone, Serialize)]
52pub struct TableMapEvent {
53 pub table_id: u64,
55 pub database_name: String,
57 pub table_name: String,
59 pub column_count: u64,
61 pub column_types: Vec<u8>,
63 #[serde(skip_serializing_if = "Vec::is_empty")]
65 pub column_metadata: Vec<u8>,
66 #[serde(skip_serializing_if = "Vec::is_empty")]
68 pub null_bitmap: Vec<u8>,
69}
70
71impl TableMapEvent {
72 pub fn parse(data: &[u8]) -> Option<Self> {
74 if data.len() < 10 {
75 return None;
76 }
77
78 let table_id = LittleEndian::read_u32(&data[0..]) as u64
80 | ((data[4] as u64) << 32)
81 | ((data[5] as u64) << 40);
82
83 let mut offset = 8;
85
86 if offset >= data.len() {
88 return None;
89 }
90 let db_len = data[offset] as usize;
91 offset += 1;
92 if offset + db_len + 1 > data.len() {
93 return None;
94 }
95 let database_name = std::str::from_utf8(&data[offset..offset + db_len])
96 .unwrap_or("")
97 .to_string();
98 offset += db_len + 1; if offset >= data.len() {
102 return None;
103 }
104 let tbl_len = data[offset] as usize;
105 offset += 1;
106 if offset + tbl_len + 1 > data.len() {
107 return None;
108 }
109 let table_name = std::str::from_utf8(&data[offset..offset + tbl_len])
110 .unwrap_or("")
111 .to_string();
112 offset += tbl_len + 1; if offset >= data.len() {
116 return None;
117 }
118 let (column_count, bytes_read) = read_lenenc_int(&data[offset..]);
119 offset += bytes_read;
120
121 let end = offset + column_count as usize;
123 if end > data.len() {
124 return None;
125 }
126 let column_types = data[offset..end].to_vec();
127 offset = end;
128
129 let column_metadata = if offset < data.len() {
131 let (meta_len, meta_bytes_read) = read_lenenc_int(&data[offset..]);
132 offset += meta_bytes_read;
133 let meta_end = offset + meta_len as usize;
134 if meta_end <= data.len() {
135 let meta = data[offset..meta_end].to_vec();
136 offset = meta_end;
137 meta
138 } else {
139 Vec::new()
140 }
141 } else {
142 Vec::new()
143 };
144
145 let null_bitmap_len = (column_count as usize).div_ceil(8);
147 let null_bitmap = if offset + null_bitmap_len <= data.len() {
148 let bm = data[offset..offset + null_bitmap_len].to_vec();
149 #[allow(unused_assignments)]
150 {
151 offset += null_bitmap_len;
152 }
153 bm
154 } else {
155 Vec::new()
156 };
157
158 Some(TableMapEvent {
159 table_id,
160 database_name,
161 table_name,
162 column_count,
163 column_types,
164 column_metadata,
165 null_bitmap,
166 })
167 }
168}
169
170#[derive(Debug, Clone, Serialize)]
174pub struct RowsEvent {
175 pub table_id: u64,
177 pub event_type: BinlogEventType,
179 pub flags: u16,
181 pub column_count: u64,
183 pub row_count: usize,
185 #[serde(skip)]
187 pub columns_present: Vec<u8>,
188 #[serde(skip)]
190 pub row_data: Vec<u8>,
191}
192
193impl RowsEvent {
194 pub fn parse(data: &[u8], type_code: u8) -> Option<Self> {
198 if data.len() < 10 {
199 return None;
200 }
201
202 let event_type = BinlogEventType::from_u8(type_code);
203
204 let table_id = LittleEndian::read_u32(&data[0..]) as u64
206 | ((data[4] as u64) << 32)
207 | ((data[5] as u64) << 40);
208
209 let flags = LittleEndian::read_u16(&data[6..]);
211
212 let extra_len = LittleEndian::read_u16(&data[8..]) as usize;
214 let mut offset = 10 + extra_len.saturating_sub(2); if offset >= data.len() {
218 return None;
219 }
220 let (column_count, bytes_read) = read_lenenc_int(&data[offset..]);
221 offset += bytes_read;
222
223 let bitmap_len = (column_count as usize).div_ceil(8);
225 let columns_present = if offset + bitmap_len <= data.len() {
226 let bm = data[offset..offset + bitmap_len].to_vec();
227 offset += bitmap_len;
228 bm
229 } else {
230 offset += bitmap_len;
231 vec![0xFF; bitmap_len]
232 };
233 if type_code == 31 {
234 offset += bitmap_len; }
236
237 let row_data = if offset < data.len() {
239 data[offset..].to_vec()
240 } else {
241 Vec::new()
242 };
243 let row_count = if !row_data.is_empty() { 1 } else { 0 };
244
245 Some(RowsEvent {
246 table_id,
247 event_type,
248 flags,
249 column_count,
250 row_count,
251 columns_present,
252 row_data,
253 })
254 }
255}
256
257pub fn read_lenenc_int(data: &[u8]) -> (u64, usize) {
261 if data.is_empty() {
262 return (0, 0);
263 }
264 match data[0] {
265 0..=250 => (data[0] as u64, 1),
266 252 => {
267 if data.len() < 3 {
268 return (0, 1);
269 }
270 (LittleEndian::read_u16(&data[1..]) as u64, 3)
271 }
272 253 => {
273 if data.len() < 4 {
274 return (0, 1);
275 }
276 let v = data[1] as u64 | (data[2] as u64) << 8 | (data[3] as u64) << 16;
277 (v, 4)
278 }
279 254 => {
280 if data.len() < 9 {
281 return (0, 1);
282 }
283 (LittleEndian::read_u64(&data[1..]), 9)
284 }
285 _ => (0, 1), }
287}
288
289#[derive(Debug, Clone, Serialize)]
291pub struct BinlogEventSummary {
292 pub offset: u64,
294 pub event_type: String,
296 pub type_code: u8,
298 pub timestamp: u32,
300 pub server_id: u32,
302 pub event_length: u32,
304}
305
306#[derive(Debug, Clone, Serialize)]
308pub struct BinlogAnalysis {
309 pub format_description: FormatDescriptionEvent,
311 pub event_count: usize,
313 pub event_type_counts: std::collections::HashMap<String, usize>,
315 #[serde(skip_serializing_if = "Vec::is_empty")]
317 pub table_maps: Vec<TableMapEvent>,
318 #[serde(skip_serializing_if = "Vec::is_empty")]
320 pub events: Vec<BinlogEventSummary>,
321}
322
323use crate::binlog::constants::COMMON_HEADER_SIZE;
324use crate::binlog::header::{validate_binlog_magic, BinlogEventHeader, FormatDescriptionEvent};
325use std::io::{Read, Seek, SeekFrom};
326
327pub fn analyze_binlog<R: Read + Seek>(mut reader: R) -> Result<BinlogAnalysis, crate::IdbError> {
331 let mut magic = [0u8; 4];
333 reader
334 .read_exact(&mut magic)
335 .map_err(|e| crate::IdbError::Io(format!("Failed to read binlog magic: {e}")))?;
336
337 if !validate_binlog_magic(&magic) {
338 return Err(crate::IdbError::Parse(
339 "Not a valid MySQL binary log file (bad magic)".to_string(),
340 ));
341 }
342
343 let file_size = reader
344 .seek(SeekFrom::End(0))
345 .map_err(|e| crate::IdbError::Io(format!("Failed to seek: {e}")))?;
346 reader
347 .seek(SeekFrom::Start(4))
348 .map_err(|e| crate::IdbError::Io(format!("Failed to seek: {e}")))?;
349
350 let mut events = Vec::new();
351 let mut event_type_counts = std::collections::HashMap::new();
352 let mut table_maps = Vec::new();
353 let mut format_desc = None;
354
355 let mut position = 4u64;
356 let mut header_buf = vec![0u8; COMMON_HEADER_SIZE];
357
358 while position + COMMON_HEADER_SIZE as u64 <= file_size {
359 if reader.read_exact(&mut header_buf).is_err() {
360 break;
361 }
362
363 let hdr = match BinlogEventHeader::parse(&header_buf) {
364 Some(h) => h,
365 None => break,
366 };
367
368 if hdr.event_length < COMMON_HEADER_SIZE as u32 {
369 break;
370 }
371
372 let data_len = hdr.event_length as usize - COMMON_HEADER_SIZE;
373 let mut event_data = vec![0u8; data_len];
374 if reader.read_exact(&mut event_data).is_err() {
375 break;
376 }
377
378 let event_type = BinlogEventType::from_u8(hdr.type_code);
379
380 if hdr.type_code == 15 && format_desc.is_none() {
382 format_desc = FormatDescriptionEvent::parse(&event_data);
383 } else if hdr.type_code == 19 {
384 if let Some(tme) = TableMapEvent::parse(&event_data) {
385 table_maps.push(tme);
386 }
387 }
388
389 *event_type_counts
390 .entry(event_type.name().to_string())
391 .or_insert(0) += 1;
392
393 events.push(BinlogEventSummary {
394 offset: position,
395 event_type: event_type.name().to_string(),
396 type_code: hdr.type_code,
397 timestamp: hdr.timestamp,
398 server_id: hdr.server_id,
399 event_length: hdr.event_length,
400 });
401
402 position = if hdr.next_position > 0 {
403 hdr.next_position as u64
404 } else {
405 position + hdr.event_length as u64
406 };
407
408 if reader.seek(SeekFrom::Start(position)).is_err() {
410 break;
411 }
412 }
413
414 let format_description = format_desc.unwrap_or(FormatDescriptionEvent {
415 binlog_version: 0,
416 server_version: "unknown".to_string(),
417 create_timestamp: 0,
418 header_length: 19,
419 checksum_alg: 0,
420 });
421
422 Ok(BinlogAnalysis {
423 format_description,
424 event_count: events.len(),
425 event_type_counts,
426 table_maps,
427 events,
428 })
429}
430
431#[cfg(test)]
432mod tests {
433 use super::*;
434
435 #[test]
436 fn test_event_type_from_u8() {
437 assert_eq!(BinlogEventType::from_u8(2), BinlogEventType::QueryEvent);
438 assert_eq!(
439 BinlogEventType::from_u8(15),
440 BinlogEventType::FormatDescription
441 );
442 assert_eq!(BinlogEventType::from_u8(19), BinlogEventType::TableMapEvent);
443 assert_eq!(
444 BinlogEventType::from_u8(30),
445 BinlogEventType::WriteRowsEvent
446 );
447 assert_eq!(
448 BinlogEventType::from_u8(31),
449 BinlogEventType::UpdateRowsEvent
450 );
451 assert_eq!(
452 BinlogEventType::from_u8(32),
453 BinlogEventType::DeleteRowsEvent
454 );
455 assert_eq!(BinlogEventType::from_u8(255), BinlogEventType::Unknown(255));
456 }
457
458 #[test]
459 fn test_event_type_names() {
460 assert_eq!(
461 BinlogEventType::FormatDescription.name(),
462 "FORMAT_DESCRIPTION"
463 );
464 assert_eq!(BinlogEventType::TableMapEvent.name(), "TABLE_MAP");
465 assert_eq!(BinlogEventType::WriteRowsEvent.name(), "WRITE_ROWS_V2");
466 assert_eq!(BinlogEventType::GtidLogEvent.name(), "GTID");
467 }
468
469 #[test]
470 fn test_event_type_display() {
471 assert_eq!(format!("{}", BinlogEventType::QueryEvent), "QUERY");
472 assert_eq!(format!("{}", BinlogEventType::Unknown(99)), "UNKNOWN(99)");
473 }
474
475 #[test]
476 fn test_table_map_event_parse() {
477 let mut data = vec![0u8; 50];
478 LittleEndian::write_u32(&mut data[0..], 42);
480 data[4] = 0;
481 data[5] = 0;
482 LittleEndian::write_u16(&mut data[6..], 0);
484 data[8] = 4;
486 data[9..13].copy_from_slice(b"test");
487 data[13] = 0;
488 data[14] = 5;
490 data[15..20].copy_from_slice(b"users");
491 data[20] = 0;
492 data[21] = 3;
494 data[22] = 3; data[23] = 15; data[24] = 12; let tme = TableMapEvent::parse(&data).unwrap();
500 assert_eq!(tme.table_id, 42);
501 assert_eq!(tme.database_name, "test");
502 assert_eq!(tme.table_name, "users");
503 assert_eq!(tme.column_count, 3);
504 assert_eq!(tme.column_types, vec![3, 15, 12]);
505 }
506
507 #[test]
508 fn test_table_map_event_too_short() {
509 let data = vec![0u8; 5];
510 assert!(TableMapEvent::parse(&data).is_none());
511 }
512
513 #[test]
514 fn test_rows_event_parse() {
515 let mut data = vec![0u8; 30];
516 LittleEndian::write_u32(&mut data[0..], 42);
518 data[4] = 0;
519 data[5] = 0;
520 LittleEndian::write_u16(&mut data[6..], 1);
522 LittleEndian::write_u16(&mut data[8..], 2);
524 data[10] = 3;
526 data[11] = 0x07;
528 data[12] = 0x01;
530
531 let re = RowsEvent::parse(&data, 30).unwrap();
532 assert_eq!(re.table_id, 42);
533 assert_eq!(re.event_type, BinlogEventType::WriteRowsEvent);
534 assert_eq!(re.flags, 1);
535 assert_eq!(re.column_count, 3);
536 }
537
538 #[test]
539 fn test_lenenc_int() {
540 assert_eq!(read_lenenc_int(&[5]), (5, 1));
541 assert_eq!(read_lenenc_int(&[250]), (250, 1));
542 assert_eq!(read_lenenc_int(&[252, 0x01, 0x00]), (1, 3));
543 assert_eq!(read_lenenc_int(&[253, 0x01, 0x00, 0x00]), (1, 4));
544 }
545
546 #[test]
547 fn test_analyze_binlog_synthetic() {
548 use std::io::Cursor;
549
550 let mut binlog = Vec::new();
553 binlog.extend_from_slice(&[0xfe, 0x62, 0x69, 0x6e]); let fde_data_len = 100usize;
557 let fde_event_len = (COMMON_HEADER_SIZE + fde_data_len) as u32;
558
559 let mut fde_header = vec![0u8; 19];
560 LittleEndian::write_u32(&mut fde_header[0..], 1700000000); fde_header[4] = 15; LittleEndian::write_u32(&mut fde_header[5..], 1); LittleEndian::write_u32(&mut fde_header[9..], fde_event_len); LittleEndian::write_u32(&mut fde_header[13..], 4 + fde_event_len); binlog.extend_from_slice(&fde_header);
566
567 let mut fde_data = vec![0u8; fde_data_len];
568 LittleEndian::write_u16(&mut fde_data[0..], 4); let ver = b"8.0.35";
570 fde_data[2..2 + ver.len()].copy_from_slice(ver);
571 LittleEndian::write_u32(&mut fde_data[52..], 1700000000);
572 fde_data[56] = 19;
573 fde_data[95] = 1; binlog.extend_from_slice(&fde_data);
575
576 let cursor = Cursor::new(binlog);
577 let analysis = analyze_binlog(cursor).unwrap();
578
579 assert_eq!(analysis.event_count, 1);
580 assert_eq!(analysis.format_description.binlog_version, 4);
581 assert_eq!(analysis.format_description.server_version, "8.0.35");
582 assert_eq!(
583 analysis.event_type_counts.get("FORMAT_DESCRIPTION"),
584 Some(&1)
585 );
586 }
587
588 #[test]
589 fn test_analyze_binlog_bad_magic() {
590 use std::io::Cursor;
591 let data = vec![0u8; 100];
592 let cursor = Cursor::new(data);
593 assert!(analyze_binlog(cursor).is_err());
594 }
595}