use crate::{BAD_HANDLE, Instance, arg, with};
#[unsafe(no_mangle)]
pub unsafe extern "C" fn kevy_aof_frame_in(h: u32, p: *const u8, l: u32) -> i32 {
let chunk = unsafe { arg(p, l) };
with(h, BAD_HANDLE, |inst| feed(inst, chunk))
}
#[unsafe(no_mangle)]
pub extern "C" fn kevy_aof_frames_out(h: u32) -> i32 {
with(h, BAD_HANDLE, |inst| {
inst.out.clear();
std::mem::swap(&mut inst.out, &mut inst.aof_out);
inst.out.len() as i32
})
}
#[unsafe(no_mangle)]
pub extern "C" fn kevy_aof_dump(h: u32) -> i32 {
with(h, BAD_HANDLE, |inst| {
inst.aof_out.clear();
inst.aof_format = kevy_persist::AofFormat::V2;
inst.aof_out_started = true; inst.out = inst.store.dump_aof_buf();
inst.out.len() as i32
})
}
fn feed(inst: &mut Instance, chunk: &[u8]) -> i32 {
inst.aof_in_carry.extend_from_slice(chunk);
if !inst.aof_in_started {
let v1 = kevy_persist::AOF_MAGIC;
let v2 = kevy_persist::AOF2_MAGIC;
let head = &inst.aof_in_carry;
if head.len() < v2.len()
&& (v1.starts_with(head.as_slice()) || v2.starts_with(head.as_slice()))
{
return 0;
}
if head.starts_with(v2) {
inst.aof_format = kevy_persist::AofFormat::V2;
inst.aof_in_carry.drain(..v2.len());
} else {
inst.aof_format = kevy_persist::AofFormat::V1;
if inst.aof_in_carry.starts_with(v1) {
inst.aof_in_carry.drain(..v1.len());
}
}
inst.aof_out_started = true;
inst.aof_in_started = true;
}
match inst.aof_format {
kevy_persist::AofFormat::V1 => feed_v1(inst),
kevy_persist::AofFormat::V2 => feed_v2(inst),
}
}
fn feed_v1(inst: &mut Instance) -> i32 {
let mut pos = 0;
let mut applied = 0i32;
loop {
match kevy_resp::parse_command(&inst.aof_in_carry[pos..]) {
Ok(Some((args, consumed))) => {
inst.store.apply_frame(&args);
pos += consumed;
applied += 1;
}
Ok(None) => break, Err(e) => {
inst.aof_in_carry.clear();
return inst
.fail(format!("corrupt AOF frame after {applied} applied frame(s): {e:?}"));
}
}
}
inst.aof_in_carry.drain(..pos);
applied
}
fn feed_v2(inst: &mut Instance) -> i32 {
let mut pos = 0;
let mut applied = 0i32;
loop {
match kevy_persist::next_record(&inst.aof_in_carry, pos) {
kevy_persist::RecordStep::Ok { payload, consumed } => {
match kevy_resp::parse_command(payload) {
Ok(Some((args, used))) if used == payload.len() => {
inst.store.apply_frame(&args);
pos += consumed;
applied += 1;
}
_ => {
inst.aof_in_carry.clear();
return inst.fail(format!(
"corrupt AOF record (bad payload) after {applied} applied frame(s)"
));
}
}
}
kevy_persist::RecordStep::Truncated => break,
kevy_persist::RecordStep::Corrupt => {
inst.aof_in_carry.clear();
return inst.fail(format!("corrupt AOF record after {applied} applied frame(s)"));
}
}
}
inst.aof_in_carry.drain(..pos);
applied
}