Skip to main content

idb/cli/
binlog.rs

1//! CLI implementation for the `inno binlog` subcommand.
2//!
3//! Parses MySQL binary log files and displays event summaries, format
4//! description info, table map details, and row-based event statistics.
5//! With `--correlate`, maps row events to tablespace pages via B+Tree lookup.
6
7use std::collections::HashMap;
8use std::io::Write;
9
10use serde::Serialize;
11
12use crate::cli::{csv_escape, wprintln};
13use crate::IdbError;
14
15/// Options for the `inno binlog` subcommand.
16pub struct BinlogOptions {
17    /// Path to the MySQL binary log file.
18    pub file: String,
19    /// Maximum number of events to display.
20    pub limit: Option<usize>,
21    /// Filter events by type name (e.g. "TABLE_MAP", "WRITE_ROWS").
22    pub filter_type: Option<String>,
23    /// Show additional detail (column types for TABLE_MAP events).
24    pub verbose: bool,
25    /// Output in JSON format.
26    pub json: bool,
27    /// Output as CSV.
28    pub csv: bool,
29    /// Path to .ibd tablespace for page correlation.
30    pub correlate: Option<String>,
31}
32
33/// Combined analysis with correlated events for JSON output.
34#[derive(Serialize)]
35struct CorrelatedBinlogAnalysis {
36    #[serde(flatten)]
37    analysis: crate::binlog::BinlogAnalysis,
38    correlated_events: Vec<crate::binlog::CorrelatedEvent>,
39}
40
41/// Analyze a binary log file and display results.
42pub fn execute(opts: &BinlogOptions, writer: &mut dyn Write) -> Result<(), IdbError> {
43    if opts.correlate.is_some() {
44        return execute_correlated(opts, writer);
45    }
46
47    let file = std::fs::File::open(&opts.file)
48        .map_err(|e| IdbError::Io(format!("{}: {}", opts.file, e)))?;
49
50    let reader = std::io::BufReader::new(file);
51    let analysis = crate::binlog::analyze_binlog(reader)?;
52
53    if opts.json {
54        let json =
55            serde_json::to_string_pretty(&analysis).map_err(|e| IdbError::Parse(e.to_string()))?;
56        wprintln!(writer, "{}", json)?;
57        return Ok(());
58    }
59
60    if opts.csv {
61        return write_csv(&analysis, opts, writer);
62    }
63
64    write_text(&analysis, opts, writer)
65}
66
67/// Execute with page correlation: maps row events to tablespace pages.
68fn execute_correlated(opts: &BinlogOptions, writer: &mut dyn Write) -> Result<(), IdbError> {
69    let ts_path = opts.correlate.as_ref().unwrap();
70
71    // Run correlation
72    let mut binlog = crate::binlog::BinlogFile::open(&opts.file)?;
73    let mut ts = crate::cli::open_tablespace(ts_path, None, false)?;
74    let correlated = crate::binlog::correlate_events(&mut binlog, &mut ts)?;
75
76    // Also run standard analysis for event context
77    let file = std::fs::File::open(&opts.file)
78        .map_err(|e| IdbError::Io(format!("{}: {}", opts.file, e)))?;
79    let reader = std::io::BufReader::new(file);
80    let analysis = crate::binlog::analyze_binlog(reader)?;
81
82    if opts.json {
83        let combined = CorrelatedBinlogAnalysis {
84            analysis,
85            correlated_events: correlated,
86        };
87        let json =
88            serde_json::to_string_pretty(&combined).map_err(|e| IdbError::Parse(e.to_string()))?;
89        wprintln!(writer, "{}", json)?;
90        return Ok(());
91    }
92
93    // Build lookup by binlog position
94    let correlated_map: HashMap<u64, &crate::binlog::CorrelatedEvent> =
95        correlated.iter().map(|e| (e.binlog_pos, e)).collect();
96
97    if opts.csv {
98        return write_correlated_csv(&analysis, &correlated_map, opts, writer);
99    }
100
101    write_correlated_text(&analysis, &correlated_map, opts, writer)
102}
103
104/// Write text output for correlated binlog analysis.
105fn write_correlated_text(
106    analysis: &crate::binlog::BinlogAnalysis,
107    correlated: &HashMap<u64, &crate::binlog::CorrelatedEvent>,
108    opts: &BinlogOptions,
109    writer: &mut dyn Write,
110) -> Result<(), IdbError> {
111    // Format description header
112    wprintln!(writer, "Binary Log: {}", opts.file)?;
113    wprintln!(
114        writer,
115        "  Server Version: {}",
116        analysis.format_description.server_version
117    )?;
118    wprintln!(
119        writer,
120        "  Binlog Version: {}",
121        analysis.format_description.binlog_version
122    )?;
123    wprintln!(
124        writer,
125        "  Checksum Algorithm: {}",
126        analysis.format_description.checksum_alg
127    )?;
128    wprintln!(
129        writer,
130        "  Correlated Events: {} (tablespace: {})",
131        correlated.len(),
132        opts.correlate.as_deref().unwrap_or("--")
133    )?;
134    wprintln!(writer)?;
135
136    // Event type summary
137    wprintln!(
138        writer,
139        "Event Type Summary ({} total):",
140        analysis.event_count
141    )?;
142    let mut type_counts: Vec<_> = analysis.event_type_counts.iter().collect();
143    type_counts.sort_by(|a, b| b.1.cmp(a.1));
144    for (name, count) in &type_counts {
145        wprintln!(writer, "  {:<30} {:>6}", name, count)?;
146    }
147    wprintln!(writer)?;
148
149    // Table maps
150    if !analysis.table_maps.is_empty() {
151        wprintln!(writer, "Table Maps ({}):", analysis.table_maps.len())?;
152        for tm in &analysis.table_maps {
153            wprintln!(
154                writer,
155                "  table_id={} {}.{} ({} columns)",
156                tm.table_id,
157                tm.database_name,
158                tm.table_name,
159                tm.column_count
160            )?;
161            if opts.verbose {
162                wprintln!(writer, "    Column types: {:?}", &tm.column_types)?;
163            }
164        }
165        wprintln!(writer)?;
166    }
167
168    // Event listing with page correlation columns
169    let events = filter_events(&analysis.events, opts);
170    let limit = opts.limit.unwrap_or(events.len());
171    let display_events = &events[..limit.min(events.len())];
172
173    if !display_events.is_empty() {
174        wprintln!(
175            writer,
176            "{:<12} {:<30} {:<10} {:<12} {:<8} {}",
177            "Position",
178            "Type",
179            "Size",
180            "Timestamp",
181            "Page",
182            "PK"
183        )?;
184        wprintln!(writer, "{}", "-".repeat(90))?;
185        for evt in display_events {
186            if let Some(ce) = correlated.get(&evt.offset) {
187                let pk_display = if ce.pk_values.is_empty() {
188                    "--".to_string()
189                } else {
190                    format!("({})", ce.pk_values.join(", "))
191                };
192                wprintln!(
193                    writer,
194                    "{:<12} {:<30} {:<10} {:<12} {:<8} {}",
195                    evt.offset,
196                    evt.event_type,
197                    evt.event_length,
198                    evt.timestamp,
199                    ce.page_no,
200                    pk_display
201                )?;
202            } else {
203                wprintln!(
204                    writer,
205                    "{:<12} {:<30} {:<10} {:<12} {:<8} {}",
206                    evt.offset,
207                    evt.event_type,
208                    evt.event_length,
209                    evt.timestamp,
210                    "--",
211                    "--"
212                )?;
213            }
214        }
215    }
216
217    if events.len() > limit {
218        wprintln!(
219            writer,
220            "\n... {} more events (use --limit to show more)",
221            events.len() - limit
222        )?;
223    }
224
225    Ok(())
226}
227
228/// Write CSV output for correlated binlog events.
229fn write_correlated_csv(
230    analysis: &crate::binlog::BinlogAnalysis,
231    correlated: &HashMap<u64, &crate::binlog::CorrelatedEvent>,
232    opts: &BinlogOptions,
233    writer: &mut dyn Write,
234) -> Result<(), IdbError> {
235    wprintln!(
236        writer,
237        "position,type,size,timestamp,server_id,page_no,space_id,pk_values"
238    )?;
239
240    let events = filter_events(&analysis.events, opts);
241    let limit = opts.limit.unwrap_or(events.len());
242    let display_events = &events[..limit.min(events.len())];
243
244    for evt in display_events {
245        if let Some(ce) = correlated.get(&evt.offset) {
246            let pk_str = ce.pk_values.join(";");
247            wprintln!(
248                writer,
249                "{},{},{},{},{},{},{},{}",
250                evt.offset,
251                csv_escape(&evt.event_type),
252                evt.event_length,
253                evt.timestamp,
254                evt.server_id,
255                ce.page_no,
256                ce.space_id,
257                csv_escape(&pk_str)
258            )?;
259        } else {
260            wprintln!(
261                writer,
262                "{},{},{},{},{},,,",
263                evt.offset,
264                csv_escape(&evt.event_type),
265                evt.event_length,
266                evt.timestamp,
267                evt.server_id
268            )?;
269        }
270    }
271
272    Ok(())
273}
274
275/// Write text output for binlog analysis.
276fn write_text(
277    analysis: &crate::binlog::BinlogAnalysis,
278    opts: &BinlogOptions,
279    writer: &mut dyn Write,
280) -> Result<(), IdbError> {
281    // Format description header
282    wprintln!(writer, "Binary Log: {}", opts.file)?;
283    wprintln!(
284        writer,
285        "  Server Version: {}",
286        analysis.format_description.server_version
287    )?;
288    wprintln!(
289        writer,
290        "  Binlog Version: {}",
291        analysis.format_description.binlog_version
292    )?;
293    wprintln!(
294        writer,
295        "  Checksum Algorithm: {}",
296        analysis.format_description.checksum_alg
297    )?;
298    wprintln!(writer)?;
299
300    // Event type summary
301    wprintln!(
302        writer,
303        "Event Type Summary ({} total):",
304        analysis.event_count
305    )?;
306    let mut type_counts: Vec<_> = analysis.event_type_counts.iter().collect();
307    type_counts.sort_by(|a, b| b.1.cmp(a.1));
308    for (name, count) in &type_counts {
309        wprintln!(writer, "  {:<30} {:>6}", name, count)?;
310    }
311    wprintln!(writer)?;
312
313    // Table maps
314    if !analysis.table_maps.is_empty() {
315        wprintln!(writer, "Table Maps ({}):", analysis.table_maps.len())?;
316        for tm in &analysis.table_maps {
317            wprintln!(
318                writer,
319                "  table_id={} {}.{} ({} columns)",
320                tm.table_id,
321                tm.database_name,
322                tm.table_name,
323                tm.column_count
324            )?;
325            if opts.verbose {
326                wprintln!(writer, "    Column types: {:?}", &tm.column_types)?;
327            }
328        }
329        wprintln!(writer)?;
330    }
331
332    // Event listing
333    let events = filter_events(&analysis.events, opts);
334    let limit = opts.limit.unwrap_or(events.len());
335    let display_events = &events[..limit.min(events.len())];
336
337    if !display_events.is_empty() {
338        wprintln!(
339            writer,
340            "{:<12} {:<30} {:<10} {:<12}",
341            "Position",
342            "Type",
343            "Size",
344            "Timestamp"
345        )?;
346        wprintln!(writer, "{}", "-".repeat(66))?;
347        for evt in display_events {
348            wprintln!(
349                writer,
350                "{:<12} {:<30} {:<10} {:<12}",
351                evt.offset,
352                evt.event_type,
353                evt.event_length,
354                evt.timestamp
355            )?;
356        }
357    }
358
359    if events.len() > limit {
360        wprintln!(
361            writer,
362            "\n... {} more events (use --limit to show more)",
363            events.len() - limit
364        )?;
365    }
366
367    Ok(())
368}
369
370/// Write CSV output for binlog events.
371fn write_csv(
372    analysis: &crate::binlog::BinlogAnalysis,
373    opts: &BinlogOptions,
374    writer: &mut dyn Write,
375) -> Result<(), IdbError> {
376    wprintln!(writer, "position,type,size,timestamp,server_id")?;
377
378    let events = filter_events(&analysis.events, opts);
379    let limit = opts.limit.unwrap_or(events.len());
380    let display_events = &events[..limit.min(events.len())];
381
382    for evt in display_events {
383        wprintln!(
384            writer,
385            "{},{},{},{},{}",
386            evt.offset,
387            csv_escape(&evt.event_type),
388            evt.event_length,
389            evt.timestamp,
390            evt.server_id
391        )?;
392    }
393
394    Ok(())
395}
396
397/// Filter events by type name if a filter is provided.
398fn filter_events<'a>(
399    events: &'a [crate::binlog::BinlogEventSummary],
400    opts: &BinlogOptions,
401) -> Vec<&'a crate::binlog::BinlogEventSummary> {
402    match &opts.filter_type {
403        Some(filter) => {
404            let filter_upper = filter.to_uppercase();
405            events
406                .iter()
407                .filter(|e| e.event_type.to_uppercase().contains(&filter_upper))
408                .collect()
409        }
410        None => events.iter().collect(),
411    }
412}
413
414#[cfg(test)]
415mod tests {
416    use super::*;
417    use crate::binlog::{BinlogAnalysis, BinlogEventSummary, FormatDescriptionEvent};
418
419    fn sample_analysis() -> BinlogAnalysis {
420        let mut event_type_counts = HashMap::new();
421        event_type_counts.insert("QUERY_EVENT".to_string(), 5);
422        event_type_counts.insert("TABLE_MAP_EVENT".to_string(), 2);
423
424        BinlogAnalysis {
425            format_description: FormatDescriptionEvent {
426                binlog_version: 4,
427                server_version: "8.0.35".to_string(),
428                create_timestamp: 0,
429                header_length: 19,
430                checksum_alg: 1,
431            },
432            event_count: 7,
433            event_type_counts,
434            table_maps: Vec::new(),
435            events: vec![
436                BinlogEventSummary {
437                    offset: 4,
438                    event_type: "FORMAT_DESCRIPTION_EVENT".to_string(),
439                    type_code: 15,
440                    event_length: 119,
441                    timestamp: 1700000000,
442                    server_id: 1,
443                },
444                BinlogEventSummary {
445                    offset: 123,
446                    event_type: "QUERY_EVENT".to_string(),
447                    type_code: 2,
448                    event_length: 50,
449                    timestamp: 1700000001,
450                    server_id: 1,
451                },
452            ],
453        }
454    }
455
456    #[test]
457    fn test_write_text_output() {
458        let analysis = sample_analysis();
459        let opts = BinlogOptions {
460            file: "test-bin.000001".to_string(),
461            limit: None,
462            filter_type: None,
463            verbose: false,
464            json: false,
465            csv: false,
466            correlate: None,
467        };
468
469        let mut buf = Vec::new();
470        write_text(&analysis, &opts, &mut buf).unwrap();
471        let output = String::from_utf8(buf).unwrap();
472        assert!(output.contains("Binary Log: test-bin.000001"));
473        assert!(output.contains("Server Version: 8.0.35"));
474        assert!(output.contains("7 total"));
475        assert!(output.contains("FORMAT_DESCRIPTION_EVENT"));
476    }
477
478    #[test]
479    fn test_write_csv_output() {
480        let analysis = sample_analysis();
481        let opts = BinlogOptions {
482            file: "test-bin.000001".to_string(),
483            limit: None,
484            filter_type: None,
485            verbose: false,
486            json: false,
487            csv: true,
488            correlate: None,
489        };
490
491        let mut buf = Vec::new();
492        write_csv(&analysis, &opts, &mut buf).unwrap();
493        let output = String::from_utf8(buf).unwrap();
494        assert!(output.starts_with("position,type,size,timestamp,server_id"));
495        assert!(output.contains("FORMAT_DESCRIPTION_EVENT"));
496    }
497
498    #[test]
499    fn test_filter_events_by_type() {
500        let analysis = sample_analysis();
501        let opts = BinlogOptions {
502            file: "test".to_string(),
503            limit: None,
504            filter_type: Some("query".to_string()),
505            verbose: false,
506            json: false,
507            csv: false,
508            correlate: None,
509        };
510
511        let filtered = filter_events(&analysis.events, &opts);
512        assert_eq!(filtered.len(), 1);
513        assert_eq!(filtered[0].event_type, "QUERY_EVENT");
514    }
515
516    #[test]
517    fn test_json_output() {
518        let analysis = sample_analysis();
519        let json = serde_json::to_string_pretty(&analysis).unwrap();
520        let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
521        assert_eq!(parsed["event_count"], 7);
522    }
523
524    #[test]
525    fn test_correlated_text_output() {
526        let analysis = sample_analysis();
527        let ce = crate::binlog::CorrelatedEvent {
528            binlog_pos: 123,
529            event_type: crate::binlog::RowEventType::Insert,
530            database: "test".to_string(),
531            table: "users".to_string(),
532            page_no: 42,
533            space_id: 5,
534            page_lsn: 999,
535            pk_values: vec!["1".to_string(), "alice".to_string()],
536            timestamp: 1700000001,
537        };
538        let correlated: HashMap<u64, &crate::binlog::CorrelatedEvent> =
539            [(123, &ce)].into_iter().collect();
540        let opts = BinlogOptions {
541            file: "test-bin.000001".to_string(),
542            limit: None,
543            filter_type: None,
544            verbose: false,
545            json: false,
546            csv: false,
547            correlate: Some("/tmp/users.ibd".to_string()),
548        };
549
550        let mut buf = Vec::new();
551        write_correlated_text(&analysis, &correlated, &opts, &mut buf).unwrap();
552        let output = String::from_utf8(buf).unwrap();
553        assert!(output.contains("Correlated Events: 1"));
554        assert!(output.contains("Page"));
555        assert!(output.contains("PK"));
556        // The correlated row event should show page 42
557        assert!(output.contains("42"));
558        assert!(output.contains("(1, alice)"));
559        // Non-correlated event should show --
560        let lines: Vec<&str> = output.lines().collect();
561        let fde_line = lines
562            .iter()
563            .find(|l| l.contains("FORMAT_DESCRIPTION"))
564            .unwrap();
565        assert!(fde_line.contains("--"));
566    }
567
568    #[test]
569    fn test_correlated_csv_output() {
570        let analysis = sample_analysis();
571        let ce = crate::binlog::CorrelatedEvent {
572            binlog_pos: 123,
573            event_type: crate::binlog::RowEventType::Insert,
574            database: "test".to_string(),
575            table: "users".to_string(),
576            page_no: 42,
577            space_id: 5,
578            page_lsn: 999,
579            pk_values: vec!["1".to_string()],
580            timestamp: 1700000001,
581        };
582        let correlated: HashMap<u64, &crate::binlog::CorrelatedEvent> =
583            [(123, &ce)].into_iter().collect();
584        let opts = BinlogOptions {
585            file: "test-bin.000001".to_string(),
586            limit: None,
587            filter_type: None,
588            verbose: false,
589            json: false,
590            csv: true,
591            correlate: Some("/tmp/users.ibd".to_string()),
592        };
593
594        let mut buf = Vec::new();
595        write_correlated_csv(&analysis, &correlated, &opts, &mut buf).unwrap();
596        let output = String::from_utf8(buf).unwrap();
597        assert!(
598            output.starts_with("position,type,size,timestamp,server_id,page_no,space_id,pk_values")
599        );
600        // Correlated event should have page_no=42, space_id=5
601        assert!(output.contains(",42,5,"));
602        // Non-correlated event should have empty page/space columns
603        let fde_line = output
604            .lines()
605            .find(|l| l.contains("FORMAT_DESCRIPTION"))
606            .unwrap();
607        assert!(fde_line.ends_with(",,,"));
608    }
609
610    #[test]
611    fn test_correlated_json_output() {
612        let analysis = sample_analysis();
613        let correlated = vec![crate::binlog::CorrelatedEvent {
614            binlog_pos: 123,
615            event_type: crate::binlog::RowEventType::Insert,
616            database: "test".to_string(),
617            table: "users".to_string(),
618            page_no: 42,
619            space_id: 5,
620            page_lsn: 999,
621            pk_values: vec!["1".to_string()],
622            timestamp: 1700000001,
623        }];
624
625        let combined = CorrelatedBinlogAnalysis {
626            analysis,
627            correlated_events: correlated,
628        };
629        let json = serde_json::to_string_pretty(&combined).unwrap();
630        let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
631        assert_eq!(parsed["event_count"], 7);
632        assert_eq!(parsed["correlated_events"][0]["page_no"], 42);
633        assert_eq!(parsed["correlated_events"][0]["space_id"], 5);
634        assert_eq!(parsed["correlated_events"][0]["event_type"], "Insert");
635    }
636}