use std::io::{self, Write};
use kevy_resp::ArgvView;
use crate::crc32c::crc32c;
use crate::rewrite_fmt::write_multibulk;
pub const AOF2_MAGIC: &[u8; 9] = b"KEVYAOF2\n";
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AofFormat {
V1,
V2,
}
pub(crate) const RECORD_HEADER: usize = 8;
pub(crate) const MAX_RECORD: u32 = 1 << 30;
pub(crate) fn write_record<W: Write>(w: &mut W, payload: &[u8]) -> io::Result<()> {
w.write_all(&(payload.len() as u32).to_le_bytes())?;
w.write_all(&crc32c(payload).to_le_bytes())?;
w.write_all(payload)
}
pub fn write_record_multibulk<W: Write, A: ArgvView + ?Sized>(
w: &mut W,
args: &A,
scratch: &mut Vec<u8>,
) -> io::Result<()> {
scratch.clear();
write_multibulk(scratch, args)?;
write_record(w, scratch)
}
pub enum RecordStep<'a> {
Ok {
payload: &'a [u8],
consumed: usize,
},
Truncated,
Corrupt,
}
pub fn next_record(buf: &[u8], pos: usize) -> RecordStep<'_> {
let rest = &buf[pos..];
if rest.is_empty() {
return RecordStep::Truncated; }
if rest.len() < RECORD_HEADER {
return RecordStep::Truncated;
}
let len = u32::from_le_bytes(rest[..4].try_into().unwrap());
if len == 0 || len > MAX_RECORD {
return RecordStep::Corrupt;
}
let crc = u32::from_le_bytes(rest[4..8].try_into().unwrap());
let Some(payload) = rest.get(RECORD_HEADER..RECORD_HEADER + len as usize) else {
return RecordStep::Truncated;
};
if crc32c(payload) != crc {
return RecordStep::Corrupt;
}
RecordStep::Ok { payload, consumed: RECORD_HEADER + len as usize }
}
pub(crate) fn resync_scan(buf: &[u8], from: usize) -> Option<usize> {
let mut q = from;
while q < buf.len() {
match next_record(buf, q) {
RecordStep::Ok { payload, .. } => {
if matches!(
kevy_resp::parse_command(payload),
Ok(Some((_, used))) if used == payload.len()
) {
return Some(q);
}
q += 1;
}
RecordStep::Truncated | RecordStep::Corrupt => q += 1,
}
}
None
}