Skip to main content

idb/binlog/
events.rs

1//! Binlog row-based event parsing and analysis.
2//!
3//! Provides parsing for TABLE_MAP and row-based events (WRITE/UPDATE/DELETE
4//! ROWS v2), plus top-level binlog file analysis via [`analyze_binlog`].
5
6use byteorder::{ByteOrder, LittleEndian};
7use serde::Serialize;
8
9use super::event::BinlogEventType;
10
11/// Parsed TABLE_MAP event (type 19).
12///
13/// Maps a table_id to a database.table name and column type information.
14/// This event precedes row-based events to provide schema context.
15///
16/// # Examples
17///
18/// ```
19/// use idb::binlog::events::TableMapEvent;
20/// use byteorder::{LittleEndian, ByteOrder};
21///
22/// let mut data = vec![0u8; 50];
23/// // table_id (6 bytes LE)
24/// LittleEndian::write_u32(&mut data[0..], 42);
25/// data[4] = 0; data[5] = 0;
26/// // flags (2 bytes)
27/// LittleEndian::write_u16(&mut data[6..], 0);
28/// // database name length (1 byte) + name + NUL
29/// data[8] = 4; // "test"
30/// data[9..13].copy_from_slice(b"test");
31/// data[13] = 0; // NUL
32/// // table name length (1 byte) + name + NUL
33/// data[14] = 5; // "users"
34/// data[15..20].copy_from_slice(b"users");
35/// data[20] = 0; // NUL
36/// // column_count (packed integer)
37/// data[21] = 3;
38/// // column_types
39/// data[22] = 3; // LONG
40/// data[23] = 15; // VARCHAR
41/// data[24] = 12; // DATETIME
42///
43/// let tme = TableMapEvent::parse(&data).unwrap();
44/// assert_eq!(tme.table_id, 42);
45/// assert_eq!(tme.database_name, "test");
46/// assert_eq!(tme.table_name, "users");
47/// assert_eq!(tme.column_count, 3);
48/// assert_eq!(tme.column_types, vec![3, 15, 12]);
49/// // column_metadata and null_bitmap may be empty if not present in data
50/// ```
51#[derive(Debug, Clone, Serialize)]
52pub struct TableMapEvent {
53    /// Internal table ID.
54    pub table_id: u64,
55    /// Database (schema) name.
56    pub database_name: String,
57    /// Table name.
58    pub table_name: String,
59    /// Number of columns.
60    pub column_count: u64,
61    /// Column type codes (MySQL protocol type codes, e.g. 3=LONG, 15=VARCHAR).
62    pub column_types: Vec<u8>,
63    /// Per-column type metadata (variable-length, type-dependent).
64    #[serde(skip_serializing_if = "Vec::is_empty")]
65    pub column_metadata: Vec<u8>,
66    /// Column nullability bitmap (bit N set = column N is nullable).
67    #[serde(skip_serializing_if = "Vec::is_empty")]
68    pub null_bitmap: Vec<u8>,
69}
70
71impl TableMapEvent {
72    /// Parse a TABLE_MAP event from the event data (after the common header).
73    pub fn parse(data: &[u8]) -> Option<Self> {
74        if data.len() < 10 {
75            return None;
76        }
77
78        // table_id: 6 bytes LE
79        let table_id = LittleEndian::read_u32(&data[0..]) as u64
80            | ((data[4] as u64) << 32)
81            | ((data[5] as u64) << 40);
82
83        // flags: 2 bytes (skip)
84        let mut offset = 8;
85
86        // database name: 1-byte length + string + NUL
87        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; // skip NUL
99
100        // table name: 1-byte length + string + NUL
101        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; // skip NUL
113
114        // column_count: packed integer (lenenc)
115        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        // column_types: column_count bytes
122        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        // column_metadata: lenenc-prefixed metadata bytes
130        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        // null_bitmap: ceil(column_count / 8) bytes
146        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/// Parsed row-based event summary (types 30-32).
171///
172/// Contains metadata about row changes plus raw row data for correlation.
173#[derive(Debug, Clone, Serialize)]
174pub struct RowsEvent {
175    /// Internal table ID (matches TABLE_MAP event).
176    pub table_id: u64,
177    /// Event type.
178    pub event_type: BinlogEventType,
179    /// Event flags.
180    pub flags: u16,
181    /// Number of columns involved.
182    pub column_count: u64,
183    /// Approximate row count (estimated from data size).
184    pub row_count: usize,
185    /// Columns-present bitmap (before image). Bit N set = column N in row data.
186    #[serde(skip)]
187    pub columns_present: Vec<u8>,
188    /// Raw row image data (after column bitmaps, before CRC).
189    #[serde(skip)]
190    pub row_data: Vec<u8>,
191}
192
193impl RowsEvent {
194    /// Parse a row-based event from the event data (after the common header).
195    ///
196    /// Supports WRITE_ROWS_V2 (30), UPDATE_ROWS_V2 (31), DELETE_ROWS_V2 (32).
197    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        // table_id: 6 bytes LE
205        let table_id = LittleEndian::read_u32(&data[0..]) as u64
206            | ((data[4] as u64) << 32)
207            | ((data[5] as u64) << 40);
208
209        // flags: 2 bytes
210        let flags = LittleEndian::read_u16(&data[6..]);
211
212        // extra_data_length: 2 bytes (v2 events)
213        let extra_len = LittleEndian::read_u16(&data[8..]) as usize;
214        let mut offset = 10 + extra_len.saturating_sub(2); // extra_len includes itself
215
216        // column_count: packed integer
217        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        // Columns-present bitmap (before image)
224        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; // columns_after_image for UPDATE
235        }
236
237        // Capture remaining bytes as raw row data (may include CRC-32C at end)
238        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
257/// Read a MySQL packed (length-encoded) integer.
258///
259/// Returns (value, bytes_consumed).
260pub 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), // 251 = NULL, 255 = undefined
286    }
287}
288
289/// Summary of a single binlog event for analysis output.
290#[derive(Debug, Clone, Serialize)]
291pub struct BinlogEventSummary {
292    /// File offset of the event.
293    pub offset: u64,
294    /// Event type name.
295    pub event_type: String,
296    /// Event type code.
297    pub type_code: u8,
298    /// Unix timestamp.
299    pub timestamp: u32,
300    /// Server ID.
301    pub server_id: u32,
302    /// Total event size in bytes.
303    pub event_length: u32,
304}
305
306/// Top-level analysis result for a binary log file.
307#[derive(Debug, Clone, Serialize)]
308pub struct BinlogAnalysis {
309    /// Format description from the first event.
310    pub format_description: FormatDescriptionEvent,
311    /// Total number of events.
312    pub event_count: usize,
313    /// Count of events by type name.
314    pub event_type_counts: std::collections::HashMap<String, usize>,
315    /// TABLE_MAP events found.
316    #[serde(skip_serializing_if = "Vec::is_empty")]
317    pub table_maps: Vec<TableMapEvent>,
318    /// Individual event summaries.
319    #[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
327/// Analyze a binary log file from a reader.
328///
329/// Reads all events, collecting summaries, TABLE_MAP events, and type counts.
330pub fn analyze_binlog<R: Read + Seek>(mut reader: R) -> Result<BinlogAnalysis, crate::IdbError> {
331    // Validate magic
332    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        // Parse specific event types
381        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        // Seek to next event position (in case of padding or checksum)
409        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        // table_id = 42
479        LittleEndian::write_u32(&mut data[0..], 42);
480        data[4] = 0;
481        data[5] = 0;
482        // flags
483        LittleEndian::write_u16(&mut data[6..], 0);
484        // db name: "test"
485        data[8] = 4;
486        data[9..13].copy_from_slice(b"test");
487        data[13] = 0;
488        // table name: "users"
489        data[14] = 5;
490        data[15..20].copy_from_slice(b"users");
491        data[20] = 0;
492        // column_count = 3
493        data[21] = 3;
494        // column types
495        data[22] = 3; // LONG
496        data[23] = 15; // VARCHAR
497        data[24] = 12; // DATETIME
498
499        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        // table_id = 42
517        LittleEndian::write_u32(&mut data[0..], 42);
518        data[4] = 0;
519        data[5] = 0;
520        // flags
521        LittleEndian::write_u16(&mut data[6..], 1);
522        // extra_data_length = 2 (minimum, self-inclusive)
523        LittleEndian::write_u16(&mut data[8..], 2);
524        // column_count = 3
525        data[10] = 3;
526        // bitmap (1 byte for 3 columns)
527        data[11] = 0x07;
528        // some row data
529        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        // Build a minimal synthetic binlog:
551        // 4-byte magic + FDE event (19-byte header + FDE data)
552        let mut binlog = Vec::new();
553        binlog.extend_from_slice(&[0xfe, 0x62, 0x69, 0x6e]); // magic
554
555        // Build FDE event
556        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); // timestamp
561        fde_header[4] = 15; // FORMAT_DESCRIPTION_EVENT
562        LittleEndian::write_u32(&mut fde_header[5..], 1); // server_id
563        LittleEndian::write_u32(&mut fde_header[9..], fde_event_len); // event_length
564        LittleEndian::write_u32(&mut fde_header[13..], 4 + fde_event_len); // next_position
565        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); // binlog_version
569        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; // checksum_alg = CRC32
574        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}