use serde::Serialize;
use std::collections::HashMap;
use std::io::{Read, Seek};
use crate::innodb::log::{
compute_record_lsn, parse_mlog_records, LogBlockHeader, LogFile, LOG_FILE_HDR_BLOCKS,
};
use crate::innodb::page::FilHeader;
use crate::innodb::page_types::PageType;
use crate::innodb::undo::{parse_undo_records, UndoRecordType};
use crate::IdbError;
#[derive(Debug, Clone, Serialize, PartialEq, Eq)]
pub enum TimelineSource {
RedoLog,
UndoLog,
Binlog,
}
impl std::fmt::Display for TimelineSource {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
TimelineSource::RedoLog => write!(f, "REDO"),
TimelineSource::UndoLog => write!(f, "UNDO"),
TimelineSource::Binlog => write!(f, "BINLOG"),
}
}
}
#[derive(Debug, Clone, Serialize)]
#[serde(tag = "type")]
pub enum TimelineAction {
Redo { mlog_type: String, single_rec: bool },
Undo {
record_type: String,
trx_id: u64,
undo_no: u64,
table_id: u64,
},
Binlog {
event_type: String,
#[serde(skip_serializing_if = "Option::is_none")]
database: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
table: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
xid: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pk_values: Option<Vec<String>>,
},
}
#[derive(Debug, Clone, Serialize)]
pub struct TimelineEntry {
pub seq: u64,
pub source: TimelineSource,
#[serde(skip_serializing_if = "Option::is_none")]
pub lsn: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub timestamp: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub space_id: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub page_no: Option<u32>,
pub action: TimelineAction,
}
#[derive(Debug, Clone, Serialize)]
pub struct PageTimelineSummary {
pub space_id: u32,
pub page_no: u32,
pub redo_entries: usize,
pub undo_entries: usize,
pub binlog_entries: usize,
#[serde(skip_serializing_if = "Option::is_none")]
pub first_lsn: Option<u64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub last_lsn: Option<u64>,
}
#[derive(Debug, Clone, Serialize)]
pub struct TimelineReport {
pub redo_count: usize,
pub undo_count: usize,
pub binlog_count: usize,
pub correlated_count: usize,
pub entries: Vec<TimelineEntry>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub page_summaries: Vec<PageTimelineSummary>,
}
pub fn extract_redo_timeline(log: &mut LogFile) -> Result<Vec<TimelineEntry>, IdbError> {
let header = log.read_header()?;
let data_blocks = log.data_block_count();
let mut entries = Vec::new();
for i in 0..data_blocks {
let block_idx = LOG_FILE_HDR_BLOCKS + i;
let block_data = log.read_block(block_idx)?;
let hdr = match LogBlockHeader::parse(&block_data) {
Some(h) if h.has_data() => h,
_ => continue,
};
let records = parse_mlog_records(&block_data, &hdr);
for rec in records {
let lsn = compute_record_lsn(header.start_lsn, i, rec.block_offset);
entries.push(TimelineEntry {
seq: 0, source: TimelineSource::RedoLog,
lsn: Some(lsn),
timestamp: None,
space_id: rec.space_id,
page_no: rec.page_no,
action: TimelineAction::Redo {
mlog_type: rec.record_type.to_string(),
single_rec: rec.single_rec,
},
});
}
}
Ok(entries)
}
pub fn extract_undo_timeline(
ts: &mut crate::innodb::tablespace::Tablespace,
) -> Result<Vec<TimelineEntry>, IdbError> {
let page_count = ts.page_count();
let mut entries = Vec::new();
for pn in 0..page_count {
let page_data = ts.read_page(pn)?;
let fil = match FilHeader::parse(&page_data) {
Some(h) => h,
None => continue,
};
if fil.page_type != PageType::UndoLog {
continue;
}
let records = parse_undo_records(&page_data);
for rec in records {
let record_type_str = match rec.record_type {
UndoRecordType::InsertRec => "INSERT",
UndoRecordType::UpdExistRec => "UPDATE",
UndoRecordType::UpdDelRec => "UPDATE_DELETE",
UndoRecordType::DelMarkRec => "DELETE_MARK",
UndoRecordType::Unknown(_) => "UNKNOWN",
};
entries.push(TimelineEntry {
seq: 0,
source: TimelineSource::UndoLog,
lsn: Some(fil.lsn),
timestamp: None,
space_id: None,
page_no: Some(fil.page_number),
action: TimelineAction::Undo {
record_type: record_type_str.to_string(),
trx_id: rec.trx_id.unwrap_or(0),
undo_no: rec.undo_no,
table_id: rec.table_id,
},
});
}
}
Ok(entries)
}
pub fn extract_binlog_timeline<R: Read + Seek>(reader: R) -> Result<Vec<TimelineEntry>, IdbError> {
let analysis = crate::binlog::events::analyze_binlog(reader)?;
let mut table_map_idx: usize = 0;
let mut current_db: Option<String> = None;
let mut current_table: Option<String> = None;
let mut entries = Vec::new();
for ev in &analysis.events {
match ev.type_code {
19 => {
if table_map_idx < analysis.table_maps.len() {
let tme = &analysis.table_maps[table_map_idx];
current_db = Some(tme.database_name.clone());
current_table = Some(tme.table_name.clone());
table_map_idx += 1;
}
}
30..=32 => {
entries.push(TimelineEntry {
seq: 0,
source: TimelineSource::Binlog,
lsn: None,
timestamp: Some(ev.timestamp),
space_id: None,
page_no: None,
action: TimelineAction::Binlog {
event_type: ev.event_type.clone(),
database: current_db.clone(),
table: current_table.clone(),
xid: None,
pk_values: None,
},
});
}
2 => {
entries.push(TimelineEntry {
seq: 0,
source: TimelineSource::Binlog,
lsn: None,
timestamp: Some(ev.timestamp),
space_id: None,
page_no: None,
action: TimelineAction::Binlog {
event_type: ev.event_type.clone(),
database: None,
table: None,
xid: None,
pk_values: None,
},
});
}
_ => {}
}
}
Ok(entries)
}
pub struct BinlogExtractionResult {
pub entries: Vec<TimelineEntry>,
pub table_maps: HashMap<u64, crate::binlog::events::TableMapEvent>,
pub row_data: HashMap<usize, Vec<u8>>,
}
pub fn extract_binlog_timeline_enriched<R: Read + Seek>(
reader: R,
) -> Result<BinlogExtractionResult, IdbError> {
use crate::binlog::constants::COMMON_HEADER_SIZE;
use crate::binlog::events::{RowsEvent, TableMapEvent};
use crate::binlog::header::{validate_binlog_magic, BinlogEventHeader};
use std::io::SeekFrom;
let mut reader = reader;
let mut magic = [0u8; 4];
reader
.read_exact(&mut magic)
.map_err(|e| IdbError::Io(format!("Failed to read binlog magic: {e}")))?;
if !validate_binlog_magic(&magic) {
return Err(IdbError::Parse(
"Not a valid MySQL binary log file (bad magic)".to_string(),
));
}
let file_size = reader
.seek(SeekFrom::End(0))
.map_err(|e| IdbError::Io(format!("Failed to seek: {e}")))?;
reader
.seek(SeekFrom::Start(4))
.map_err(|e| IdbError::Io(format!("Failed to seek: {e}")))?;
let mut entries = Vec::new();
let mut table_maps: HashMap<u64, TableMapEvent> = HashMap::new();
let mut row_data_map: HashMap<usize, Vec<u8>> = HashMap::new();
let mut current_db: Option<String> = None;
let mut current_table: Option<String> = None;
let mut position = 4u64;
let mut header_buf = vec![0u8; COMMON_HEADER_SIZE];
while position + COMMON_HEADER_SIZE as u64 <= file_size {
if reader.read_exact(&mut header_buf).is_err() {
break;
}
let hdr = match BinlogEventHeader::parse(&header_buf) {
Some(h) => h,
None => break,
};
if hdr.event_length < COMMON_HEADER_SIZE as u32 {
break;
}
let data_len = hdr.event_length as usize - COMMON_HEADER_SIZE;
let mut event_data = vec![0u8; data_len];
if reader.read_exact(&mut event_data).is_err() {
break;
}
match hdr.type_code {
19 => {
if let Some(tme) = TableMapEvent::parse(&event_data) {
current_db = Some(tme.database_name.clone());
current_table = Some(tme.table_name.clone());
table_maps.insert(tme.table_id, tme);
}
}
30..=32 => {
let entry_idx = entries.len();
entries.push(TimelineEntry {
seq: 0,
source: TimelineSource::Binlog,
lsn: None,
timestamp: Some(hdr.timestamp),
space_id: None,
page_no: None,
action: TimelineAction::Binlog {
event_type: crate::binlog::event::BinlogEventType::from_u8(hdr.type_code)
.name()
.to_string(),
database: current_db.clone(),
table: current_table.clone(),
xid: None,
pk_values: None,
},
});
if let Some(rows_ev) = RowsEvent::parse(&event_data, hdr.type_code) {
if !rows_ev.row_data.is_empty() {
row_data_map.insert(entry_idx, rows_ev.row_data);
}
}
}
2 => {
entries.push(TimelineEntry {
seq: 0,
source: TimelineSource::Binlog,
lsn: None,
timestamp: Some(hdr.timestamp),
space_id: None,
page_no: None,
action: TimelineAction::Binlog {
event_type: "QUERY".to_string(),
database: None,
table: None,
xid: None,
pk_values: None,
},
});
}
_ => {}
}
position = if hdr.next_position > 0 {
hdr.next_position as u64
} else {
position + hdr.event_length as u64
};
if reader.seek(SeekFrom::Start(position)).is_err() {
break;
}
}
Ok(BinlogExtractionResult {
entries,
table_maps,
row_data: row_data_map,
})
}
pub fn correlate_binlog_pages(
entries: &mut [TimelineEntry],
ts: &mut crate::innodb::tablespace::Tablespace,
table_maps: &HashMap<u64, crate::binlog::events::TableMapEvent>,
row_data_map: &HashMap<usize, Vec<u8>>,
) -> Result<usize, IdbError> {
use crate::binlog::correlate::{
build_column_meta, convert_pk_values, extract_ddl_column_names,
};
use crate::binlog::row_image::extract_pk_from_row_image;
use crate::innodb::btree::{extract_clustered_index_info, search_btree};
let (root_page_no, index_id, pk_columns) = match extract_clustered_index_info(ts) {
Some(info) => info,
None => return Ok(0), };
let ddl_column_names = extract_ddl_column_names(ts).unwrap_or_default();
let page0 = ts.read_page(0)?;
let space_id = FilHeader::parse(&page0).map(|h| h.space_id);
let page_size = ts.page_size();
let mut correlated = 0usize;
for (entry_idx, entry) in entries.iter_mut().enumerate() {
let row_data = match row_data_map.get(&entry_idx) {
Some(data) if !data.is_empty() => data,
_ => continue,
};
let (db_name, tbl_name) = match &entry.action {
TimelineAction::Binlog {
database: Some(db),
table: Some(tbl),
..
} => (db.clone(), tbl.clone()),
_ => continue,
};
let tme = match table_maps
.values()
.find(|t| t.database_name == db_name && t.table_name == tbl_name)
{
Some(t) => t,
None => continue,
};
let columns = build_column_meta(tme, &pk_columns, &ddl_column_names);
let pk_values = match extract_pk_from_row_image(row_data, &columns) {
Some(pks) => pks,
None => continue,
};
let search_key = convert_pk_values(&pk_values);
match search_btree(
ts,
root_page_no,
index_id,
&pk_columns,
&search_key,
page_size,
) {
Ok(result) => {
entry.page_no = Some(result.leaf_page_no);
entry.space_id = space_id;
let pk_strs: Vec<String> = pk_values.iter().map(|v| v.to_string()).collect();
if let TimelineAction::Binlog {
ref mut pk_values, ..
} = entry.action
{
*pk_values = Some(pk_strs);
}
correlated += 1;
}
Err(_) => {
continue;
}
}
}
Ok(correlated)
}
pub fn merge_timeline(
mut redo: Vec<TimelineEntry>,
mut undo: Vec<TimelineEntry>,
mut binlog: Vec<TimelineEntry>,
) -> TimelineReport {
let redo_count = redo.len();
let undo_count = undo.len();
let binlog_count = binlog.len();
let mut all = Vec::with_capacity(redo_count + undo_count + binlog_count);
all.append(&mut redo);
all.append(&mut undo);
all.append(&mut binlog);
all.sort_by(|a, b| {
let lsn_cmp = a.lsn.unwrap_or(u64::MAX).cmp(&b.lsn.unwrap_or(u64::MAX));
if lsn_cmp != std::cmp::Ordering::Equal {
return lsn_cmp;
}
a.timestamp
.unwrap_or(u32::MAX)
.cmp(&b.timestamp.unwrap_or(u32::MAX))
});
for (i, entry) in all.iter_mut().enumerate() {
entry.seq = (i + 1) as u64;
}
let mut page_agg: HashMap<(u32, u32), PageTimelineSummary> = HashMap::new();
for entry in &all {
if let (Some(sid), Some(pno)) = (entry.space_id, entry.page_no) {
let summary = page_agg.entry((sid, pno)).or_insert(PageTimelineSummary {
space_id: sid,
page_no: pno,
redo_entries: 0,
undo_entries: 0,
binlog_entries: 0,
first_lsn: None,
last_lsn: None,
});
match entry.source {
TimelineSource::RedoLog => summary.redo_entries += 1,
TimelineSource::UndoLog => summary.undo_entries += 1,
TimelineSource::Binlog => summary.binlog_entries += 1,
}
if let Some(lsn) = entry.lsn {
summary.first_lsn = Some(summary.first_lsn.map_or(lsn, |v: u64| v.min(lsn)));
summary.last_lsn = Some(summary.last_lsn.map_or(lsn, |v: u64| v.max(lsn)));
}
}
}
let correlated_count = page_agg
.values()
.filter(|s| {
let sources = [s.redo_entries > 0, s.undo_entries > 0, s.binlog_entries > 0];
sources.iter().filter(|&&v| v).count() >= 2
})
.count();
let mut page_summaries: Vec<PageTimelineSummary> = page_agg.into_values().collect();
page_summaries.sort_by_key(|s| (s.space_id, s.page_no));
TimelineReport {
redo_count,
undo_count,
binlog_count,
correlated_count,
entries: all,
page_summaries,
}
}
#[cfg(not(target_arch = "wasm32"))]
pub fn build_space_table_map(datadir: &str) -> Result<HashMap<u32, String>, IdbError> {
use std::path::Path;
use crate::innodb::tablespace::Tablespace;
use crate::util::fs::find_tablespace_files;
let files = find_tablespace_files(Path::new(datadir), &["ibd"], None)?;
let mut map = HashMap::new();
for path in files {
let path_str = path.to_string_lossy().to_string();
if let Ok(mut ts) = Tablespace::open(&path_str) {
let space_id = ts
.read_page(0)
.ok()
.and_then(|p| FilHeader::parse(&p))
.map(|h| h.space_id);
if let Some(sid) = space_id {
let p = Path::new(&path_str);
let table = p.file_stem().and_then(|s| s.to_str()).unwrap_or("unknown");
let db = p
.parent()
.and_then(|d| d.file_name())
.and_then(|s| s.to_str())
.unwrap_or("");
let full = if db.is_empty() {
table.to_string()
} else {
format!("{}.{}", db, table)
};
map.insert(sid, full);
}
}
}
Ok(map)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_merge_empty() {
let report = merge_timeline(vec![], vec![], vec![]);
assert_eq!(report.redo_count, 0);
assert_eq!(report.undo_count, 0);
assert_eq!(report.binlog_count, 0);
assert_eq!(report.correlated_count, 0);
assert!(report.entries.is_empty());
assert!(report.page_summaries.is_empty());
}
#[test]
fn test_merge_sort_by_lsn() {
let redo = vec![
TimelineEntry {
seq: 0,
source: TimelineSource::RedoLog,
lsn: Some(200),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Redo {
mlog_type: "MLOG_REC_INSERT".to_string(),
single_rec: true,
},
},
TimelineEntry {
seq: 0,
source: TimelineSource::RedoLog,
lsn: Some(100),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Redo {
mlog_type: "MLOG_REC_DELETE".to_string(),
single_rec: false,
},
},
];
let undo = vec![TimelineEntry {
seq: 0,
source: TimelineSource::UndoLog,
lsn: Some(150),
timestamp: None,
space_id: None,
page_no: Some(7),
action: TimelineAction::Undo {
record_type: "INSERT".to_string(),
trx_id: 42,
undo_no: 1,
table_id: 10,
},
}];
let report = merge_timeline(redo, undo, vec![]);
assert_eq!(report.redo_count, 2);
assert_eq!(report.undo_count, 1);
assert_eq!(report.entries.len(), 3);
assert_eq!(report.entries[0].lsn, Some(100));
assert_eq!(report.entries[0].seq, 1);
assert_eq!(report.entries[1].lsn, Some(150));
assert_eq!(report.entries[1].seq, 2);
assert_eq!(report.entries[2].lsn, Some(200));
assert_eq!(report.entries[2].seq, 3);
}
#[test]
fn test_merge_page_summaries() {
let redo = vec![
TimelineEntry {
seq: 0,
source: TimelineSource::RedoLog,
lsn: Some(100),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Redo {
mlog_type: "MLOG_REC_INSERT".to_string(),
single_rec: true,
},
},
TimelineEntry {
seq: 0,
source: TimelineSource::RedoLog,
lsn: Some(200),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Redo {
mlog_type: "MLOG_REC_DELETE".to_string(),
single_rec: true,
},
},
];
let report = merge_timeline(redo, vec![], vec![]);
assert_eq!(report.page_summaries.len(), 1);
let ps = &report.page_summaries[0];
assert_eq!(ps.space_id, 5);
assert_eq!(ps.page_no, 3);
assert_eq!(ps.redo_entries, 2);
assert_eq!(ps.first_lsn, Some(100));
assert_eq!(ps.last_lsn, Some(200));
}
#[test]
fn test_merge_correlated_count() {
let redo = vec![TimelineEntry {
seq: 0,
source: TimelineSource::RedoLog,
lsn: Some(100),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Redo {
mlog_type: "MLOG_REC_INSERT".to_string(),
single_rec: true,
},
}];
let undo = vec![TimelineEntry {
seq: 0,
source: TimelineSource::UndoLog,
lsn: Some(150),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Undo {
record_type: "INSERT".to_string(),
trx_id: 1,
undo_no: 1,
table_id: 1,
},
}];
let report = merge_timeline(redo, undo, vec![]);
assert_eq!(report.correlated_count, 1);
}
#[test]
fn test_binlog_entries_sort_after_lsn() {
let redo = vec![TimelineEntry {
seq: 0,
source: TimelineSource::RedoLog,
lsn: Some(100),
timestamp: None,
space_id: Some(5),
page_no: Some(3),
action: TimelineAction::Redo {
mlog_type: "MLOG_REC_INSERT".to_string(),
single_rec: true,
},
}];
let binlog = vec![TimelineEntry {
seq: 0,
source: TimelineSource::Binlog,
lsn: None,
timestamp: Some(1700000000),
space_id: None,
page_no: None,
action: TimelineAction::Binlog {
event_type: "WRITE_ROWS_EVENT_V2".to_string(),
database: Some("test".to_string()),
table: Some("users".to_string()),
xid: None,
pk_values: None,
},
}];
let report = merge_timeline(redo, vec![], binlog);
assert_eq!(report.entries.len(), 2);
assert_eq!(report.entries[0].source, TimelineSource::RedoLog);
assert_eq!(report.entries[1].source, TimelineSource::Binlog);
}
#[test]
fn test_timeline_source_display() {
assert_eq!(TimelineSource::RedoLog.to_string(), "REDO");
assert_eq!(TimelineSource::UndoLog.to_string(), "UNDO");
assert_eq!(TimelineSource::Binlog.to_string(), "BINLOG");
}
}