use std::fs::File;
use std::io::{self, Read};
use std::path::Path;
use crate::replay_txn::{TxnMarker, txn_marker};
use kevy_resp::Argv;
pub fn replay_aof<F: FnMut(Argv)>(path: &Path, mut apply: F) -> io::Result<ReplayReport> {
if matches!(sniff_format(path)?, crate::AofFormat::V2) {
return stream_v2(path, Some(&mut apply), false);
}
let mut data = Vec::new();
match File::open(path) {
Ok(mut f) => {
f.read_to_end(&mut data)?;
}
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(ReplayReport::default()),
Err(e) => return Err(e),
}
replay_v1_slice(path, &data, &mut apply)
}
fn replay_v1_slice<F: FnMut(Argv)>(
path: &Path,
data: &[u8],
apply: &mut F,
) -> io::Result<ReplayReport> {
let total = data.len();
if total == 0 {
return Ok(ReplayReport::default());
}
let start = std::time::Instant::now();
let mut pos = if data.starts_with(crate::aof::AOF_MAGIC) {
crate::aof::AOF_MAGIC.len()
} else {
0
};
let mut replayed: u64 = 0;
let stop = loop {
if pos >= total {
break ReplayStop::Clean;
}
match kevy_resp::parse_command(&data[pos..]) {
Ok(Some((args, consumed))) => {
apply(args);
pos += consumed;
replayed += 1;
}
Ok(None) => break ReplayStop::TruncatedTail,
Err(e) => break ReplayStop::CorruptFrame(format!("{e:?}")),
}
};
let elapsed_ms = start.elapsed().as_millis();
let corrupt = matches!(stop, ReplayStop::CorruptFrame(_));
log_replay_summary(path, total, pos, replayed, &data[pos.min(total)..], stop, elapsed_ms);
Ok(ReplayReport {
commands: replayed,
bytes: total as u64,
replayed_bytes: pos as u64,
dropped_bytes: (total - pos) as u64,
corrupt,
resynced_ranges: Vec::new(),
})
}
pub fn replay_aof_resync<F: FnMut(Argv)>(path: &Path, mut apply: F) -> io::Result<ReplayReport> {
if matches!(sniff_format(path)?, crate::AofFormat::V2) {
return stream_v2(path, Some(&mut apply), true);
}
replay_aof(path, apply)
}
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct ReplayReport {
pub commands: u64,
pub bytes: u64,
pub replayed_bytes: u64,
pub dropped_bytes: u64,
pub corrupt: bool,
pub resynced_ranges: Vec<(u64, u64)>,
}
pub(crate) fn sniff_format(path: &Path) -> io::Result<crate::AofFormat> {
let mut head = [0u8; 9];
match File::open(path) {
Ok(mut f) => match f.read_exact(&mut head) {
Ok(()) if head == *crate::record::AOF2_MAGIC => Ok(crate::AofFormat::V2),
_ => Ok(crate::AofFormat::V1),
},
Err(e) if e.kind() == io::ErrorKind::NotFound => Ok(crate::AofFormat::V1),
Err(e) => Err(e),
}
}
fn stream_v2(
path: &Path,
mut apply: Option<&mut dyn FnMut(Argv)>,
resync: bool,
) -> io::Result<ReplayReport> {
use std::io::BufReader;
let file = match File::open(path) {
Ok(f) => f,
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(ReplayReport::default()),
Err(e) => return Err(e),
};
let total = file.metadata().map_or(0, |m| m.len());
let mut r = BufReader::with_capacity(256 * 1024, file);
let mut magic = [0u8; 9];
r.read_exact(&mut magic)?; let start = std::time::Instant::now();
let mut w = walk_v2(&mut r, magic.len() as u64, &mut apply)?;
let corrupt = matches!(w.stop, ReplayStop::CorruptFrame(_));
let mut ranges: Vec<(u64, u64)> = Vec::new();
if resync && corrupt {
crate::replay_resync::resync_fallback(path, &mut w, &mut apply, &mut ranges)?;
}
let elapsed_ms = start.elapsed().as_millis();
if apply.is_some() {
log_replay_summary(
path,
total as usize,
w.pos as usize,
w.replayed,
&w.preview[..w.preview_len],
w.stop,
elapsed_ms,
);
}
Ok(ReplayReport {
commands: w.replayed,
bytes: total,
replayed_bytes: w.pos,
dropped_bytes: total.saturating_sub(w.pos),
corrupt,
resynced_ranges: ranges,
})
}
pub(crate) struct V2Walk {
pub(crate) stop: ReplayStop,
pub(crate) pos: u64,
pub(crate) replayed: u64,
pub(crate) preview: [u8; 16],
pub(crate) preview_len: usize,
pub(crate) txn: Option<Vec<Argv>>,
pub(crate) txn_discarded: u64,
}
fn preview_of(bytes: &[u8], out: &mut [u8; 16]) -> usize {
let n = bytes.len().min(out.len());
out[..n].copy_from_slice(&bytes[..n]);
n
}
fn walk_v2(
r: &mut impl Read,
start_pos: u64,
apply: &mut Option<&mut dyn FnMut(Argv)>,
) -> io::Result<V2Walk> {
let mut w = V2Walk {
txn: None,
txn_discarded: 0,
stop: ReplayStop::Clean,
pos: start_pos,
replayed: 0,
preview: [0u8; 16],
preview_len: 0,
};
let mut payload: Vec<u8> = Vec::new();
w.stop = loop {
let mut header = [0u8; 8];
match read_fully(r, &mut header) {
Ok(0) => break ReplayStop::Clean,
Ok(n) if n < header.len() => break ReplayStop::TruncatedTail,
Ok(_) => {}
Err(e) => return Err(e),
}
let len = u32::from_le_bytes(header[..4].try_into().unwrap());
let crc = u32::from_le_bytes(header[4..].try_into().unwrap());
if len == 0 || len > crate::record::MAX_RECORD {
w.preview_len = preview_of(&header, &mut w.preview);
break ReplayStop::CorruptFrame(String::from("record length out of range"));
}
payload.clear();
payload.resize(len as usize, 0);
match read_fully(r, &mut payload) {
Ok(n) if n < payload.len() => break ReplayStop::TruncatedTail,
Ok(_) => {}
Err(e) => return Err(e),
}
if crate::crc32c::crc32c(&payload) != crc {
w.preview_len = preview_of(&payload, &mut w.preview);
break ReplayStop::CorruptFrame(String::from("record checksum mismatch"));
}
if !apply_record(&payload, len, apply, &mut w) {
break ReplayStop::CorruptFrame(String::from(
"checksummed record does not hold exactly one command",
));
}
};
Ok(w)
}
fn apply_record(
payload: &[u8],
len: u32,
apply: &mut Option<&mut dyn FnMut(Argv)>,
w: &mut V2Walk,
) -> bool {
match kevy_resp::parse_command(payload) {
Ok(Some((args, used))) if used == payload.len() => {
let marker = txn_marker(&args);
match marker {
Some(TxnMarker::Begin) => {
if w.txn.take().is_some() {
w.txn_discarded += 1;
}
w.txn = Some(Vec::new());
}
Some(TxnMarker::Commit) => {
if let Some(buffered) = w.txn.take()
&& let Some(f) = apply.as_deref_mut()
{
for a in buffered {
f(a);
}
}
}
None => match w.txn.as_mut() {
Some(buf) => buf.push(args),
None => {
if let Some(f) = apply.as_deref_mut() {
f(args);
}
}
},
}
w.pos += 8 + u64::from(len);
w.replayed += 1;
true
}
_ => {
w.preview_len = preview_of(payload, &mut w.preview);
false
}
}
}
fn read_fully<R: Read>(r: &mut R, buf: &mut [u8]) -> io::Result<usize> {
let mut n = 0;
while n < buf.len() {
match r.read(&mut buf[n..]) {
Ok(0) => break,
Ok(k) => n += k,
Err(e) if e.kind() == io::ErrorKind::Interrupted => continue,
Err(e) => return Err(e),
}
}
Ok(n)
}
pub(crate) fn valid_prefix_len_of_file(path: &Path, resync: bool) -> io::Result<u64> {
if matches!(sniff_format(path)?, crate::AofFormat::V2) {
return Ok(stream_v2(path, None, resync)?.replayed_bytes);
}
let mut data = Vec::new();
match File::open(path) {
Ok(mut f) => {
f.read_to_end(&mut data)?;
}
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(0),
Err(e) => return Err(e),
}
Ok(valid_prefix_len(&data) as u64)
}
fn valid_prefix_len(data: &[u8]) -> usize {
let total = data.len();
let is_v2 = data.starts_with(crate::record::AOF2_MAGIC);
let mut pos = if is_v2 || data.starts_with(crate::aof::AOF_MAGIC) {
crate::record::AOF2_MAGIC.len()
} else {
0
};
loop {
if pos >= total {
break;
}
if is_v2 {
match crate::record::next_record(data, pos) {
crate::record::RecordStep::Ok { payload, consumed } => {
match kevy_resp::parse_command(payload) {
Ok(Some((_, used))) if used == payload.len() => pos += consumed,
_ => break,
}
}
_ => break,
}
continue;
}
match kevy_resp::parse_command(&data[pos..]) {
Ok(Some((_, consumed))) => pos += consumed,
Ok(None) | Err(_) => break,
}
}
pos
}
pub(crate) enum ReplayStop {
Clean,
TruncatedTail,
CorruptFrame(String),
}
fn log_replay_summary(
path: &Path,
total: usize,
pos: usize,
replayed: u64,
remainder: &[u8],
stop: ReplayStop,
elapsed_ms: u128,
) {
let display = path.display();
let dropped = total - pos;
match stop {
ReplayStop::Clean => {
eprintln!(
"kevy: AOF {display} replayed {replayed} commands from {total} bytes \
in {elapsed_ms} ms (clean)"
);
}
ReplayStop::TruncatedTail => {
eprintln!(
"kevy: AOF {display} replayed {replayed} commands from {total} bytes \
in {elapsed_ms} ms; trailing {dropped} bytes \
were a partial frame (crash mid-append, recoverable)"
);
}
ReplayStop::CorruptFrame(err) => {
let preview = preview_bytes(remainder);
eprintln!(
"kevy WARN: AOF {display} replayed {replayed} commands in {elapsed_ms} ms \
then hit a corrupt \
frame at byte {pos}; dropping the trailing {dropped} bytes \
(quarantined before truncation). \
Preview: {preview}. Parser error: {err}. \
Most common cause: the process was killed mid-append (a torn \
frame); less commonly, non-kevy bytes got written into this \
file path (e.g. a deploy pipeline redirecting stderr here)."
);
}
}
}
fn preview_bytes(b: &[u8]) -> String {
use std::fmt::Write;
let n = b.len().min(16);
let mut hex = String::with_capacity(n * 3);
let mut ascii = String::with_capacity(n);
for &x in &b[..n] {
if !hex.is_empty() {
hex.push(' ');
}
let _ = write!(hex, "{x:02x}");
ascii.push(if (0x20..0x7f).contains(&x) { x as char } else { '.' });
}
format!("hex=[{hex}] ascii=[{ascii}]")
}