use nodedb_wal::WalRecord;
use nodedb_wal::record::{RecordType, WalRecordArgs};
use super::RedoRecord;
use crate::data::executor::core_loop::CoreLoop;
fn reconstitute_redo_records(records: &[WalRecord]) -> Vec<WalRecord> {
let mut out = Vec::new();
for record in records {
if RecordType::from_raw(record.logical_record_type()) != Some(RecordType::TransactionRedo) {
continue;
}
let redo = match RedoRecord::from_bytes(&record.payload) {
Ok(r) => r,
Err(e) => {
tracing::warn!(
lsn = record.header.lsn,
error = %e,
"skipping malformed TransactionRedo WAL record"
);
continue;
}
};
for sub in redo.ops {
match WalRecord::new(WalRecordArgs {
record_type: sub.record_type,
lsn: record.header.lsn,
tenant_id: record.header.tenant_id,
vshard_id: record.header.vshard_id,
database_id: record.header.database_id,
payload: sub.payload,
encryption_key: None,
preamble_bytes: None,
}) {
Ok(wr) => out.push(wr),
Err(e) => tracing::warn!(
lsn = record.header.lsn,
sub_record_type = sub.record_type,
error = %e,
"skipping redo sub-record that failed to reconstitute"
),
}
}
}
out
}
impl CoreLoop {
pub(crate) fn replay_transaction_redo_wal(
&mut self,
records: &[WalRecord],
num_cores: usize,
tombstones: &nodedb_wal::TombstoneSet,
) {
let reconstituted = reconstitute_redo_records(records);
if reconstituted.is_empty() {
return;
}
self.replay_vector_wal(&reconstituted, num_cores, tombstones);
self.replay_vector_extended_wal(&reconstituted, num_cores, tombstones);
self.replay_kv_wal(&reconstituted, num_cores, tombstones);
self.replay_timeseries_wal(&reconstituted, num_cores, tombstones);
self.replay_array_wal(&reconstituted, num_cores, tombstones);
self.replay_fts_wal(&reconstituted, num_cores, tombstones);
self.replay_spatial_wal(&reconstituted, num_cores, tombstones);
self.replay_document_redo(&reconstituted, num_cores, tombstones);
self.replay_graph_redo(&reconstituted, num_cores, tombstones);
self.replay_graph_node_labels_redo(&reconstituted, num_cores);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::wal::{RedoRecord, RedoSubRecord};
fn redo_wal_record(lsn: u64, tenant_id: u64, vshard_id: u32, record: &RedoRecord) -> WalRecord {
WalRecord::new(WalRecordArgs {
record_type: RecordType::TransactionRedo as u32,
lsn,
tenant_id,
vshard_id,
database_id: 0,
payload: record.to_bytes().expect("encode redo record"),
encryption_key: None,
preamble_bytes: None,
})
.expect("wal record")
}
#[test]
fn reconstitute_preserves_type_payload_and_header() {
let redo = RedoRecord {
version: 1,
ops: vec![
RedoSubRecord {
record_type: RecordType::VectorPut as u32,
payload: vec![1, 2, 3],
},
RedoSubRecord {
record_type: RecordType::SpatialPut as u32,
payload: vec![4, 5],
},
],
calvin_stamp: None,
};
let outer = redo_wal_record(77, 9, 3, &redo);
let recon = reconstitute_redo_records(std::slice::from_ref(&outer));
assert_eq!(recon.len(), 2);
assert_eq!(recon[0].logical_record_type(), RecordType::VectorPut as u32);
assert_eq!(recon[0].payload, vec![1, 2, 3]);
assert_eq!(
recon[1].logical_record_type(),
RecordType::SpatialPut as u32
);
assert_eq!(recon[1].payload, vec![4, 5]);
for r in &recon {
assert_eq!(r.header.lsn, 77);
assert_eq!(r.header.tenant_id, 9);
assert_eq!(r.header.vshard_id, 3);
}
}
#[test]
fn reconstitute_skips_non_redo_and_malformed() {
let put = WalRecord::new(WalRecordArgs {
record_type: RecordType::Put as u32,
lsn: 1,
tenant_id: 0,
vshard_id: 0,
database_id: 0,
payload: vec![9, 9, 9],
encryption_key: None,
preamble_bytes: None,
})
.expect("wal record");
let bad = WalRecord::new(WalRecordArgs {
record_type: RecordType::TransactionRedo as u32,
lsn: 2,
tenant_id: 0,
vshard_id: 0,
database_id: 0,
payload: vec![0xff, 0xff, 0xff],
encryption_key: None,
preamble_bytes: None,
})
.expect("wal record");
let recon = reconstitute_redo_records(&[put, bad]);
assert!(recon.is_empty());
}
}