1use std::collections::HashMap;
8use std::io::Write;
9
10use serde::Serialize;
11
12use crate::cli::{csv_escape, wprintln};
13use crate::IdbError;
14
15pub struct BinlogOptions {
17 pub file: String,
19 pub limit: Option<usize>,
21 pub filter_type: Option<String>,
23 pub verbose: bool,
25 pub json: bool,
27 pub csv: bool,
29 pub correlate: Option<String>,
31}
32
33#[derive(Serialize)]
35struct CorrelatedBinlogAnalysis {
36 #[serde(flatten)]
37 analysis: crate::binlog::BinlogAnalysis,
38 correlated_events: Vec<crate::binlog::CorrelatedEvent>,
39}
40
41pub 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
67fn execute_correlated(opts: &BinlogOptions, writer: &mut dyn Write) -> Result<(), IdbError> {
69 let ts_path = opts.correlate.as_ref().unwrap();
70
71 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 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 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
104fn 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 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 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 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 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
228fn 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
275fn write_text(
277 analysis: &crate::binlog::BinlogAnalysis,
278 opts: &BinlogOptions,
279 writer: &mut dyn Write,
280) -> Result<(), IdbError> {
281 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 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 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 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
370fn 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
397fn 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 assert!(output.contains("42"));
558 assert!(output.contains("(1, alice)"));
559 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 assert!(output.contains(",42,5,"));
602 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}