use std::io::{self, Read};
use kevy_resp::Argv;
use crate::replay_txn::{TxnMarker, txn_marker};
pub(crate) enum ReplayStop {
Clean,
TruncatedTail,
LengthOutranFile {
claimed: u64,
available: u64,
},
CorruptFrame(String),
}
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,
}
pub(crate) 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 outran(claimed: u32, available: usize) -> ReplayStop {
ReplayStop::LengthOutranFile { claimed: u64::from(claimed), available: available as u64 }
}
pub(crate) 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 outran(len, n),
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)
}
pub(crate) 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
}
}
}
pub(crate) 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)
}