use crate::SNAPSHOT_BUF_CAP;
use kevy_resp::ArgvView;
use crate::rewrite_chunk::write_verb_items;
use kevy_store::Value;
use std::fs::File;
use std::io::{self, BufWriter, Write};
use std::path::Path;
pub(crate) fn emit<W: Write, A: kevy_resp::ArgvView + ?Sized>(
w: &mut W,
args: &A,
fmt: crate::AofFormat,
scratch: &mut Vec<u8>,
) -> io::Result<()> {
match fmt {
crate::AofFormat::V1 => write_multibulk(w, args),
crate::AofFormat::V2 => crate::record::write_record_multibulk(w, args, scratch),
}
}
pub fn dump_aof<S: crate::SnapshotSource>(path: &Path, src: &S) -> io::Result<(u64, u64)> {
const DROP_BEHIND: u64 = 64 << 20;
let f = File::create(path)?;
let mut w = BufWriter::with_capacity(SNAPSHOT_BUF_CAP, f);
let mut scratch = Vec::new();
w.write_all(crate::record::AOF2_MAGIC)?;
let mut keys = 0u64;
let mut err: Option<io::Error> = None;
let mut cold_seqs: Vec<u32> = Vec::new();
let mut since_drop: u64 = 0;
src.for_each_entry(|key, value, ttl_ms| {
if err.is_some() {
return;
}
if let kevy_store::Value::Cold(c) = value
&& let Some((seq, _)) = c.seg_parts()
{
if !cold_seqs.contains(&seq) {
cold_seqs.push(seq);
}
return;
}
if let Err(e) =
write_value_as_commands(&mut w, key, value, ttl_ms, crate::AofFormat::V2, &mut scratch)
{
err = Some(e);
} else {
keys += 1;
since_drop = since_drop.saturating_add(scratch.len() as u64);
if since_drop >= DROP_BEHIND {
since_drop = 0;
crate::dump_cache::drop_behind_step(&mut w);
}
}
});
if let Some(e) = err {
return Err(e);
}
finish_dump(w, src, &cold_seqs, &mut scratch, keys)
}
fn finish_dump<S: crate::SnapshotSource>(
mut w: BufWriter<File>,
src: &S,
cold_seqs: &[u32],
scratch: &mut Vec<u8>,
keys: u64,
) -> io::Result<(u64, u64)> {
write_hash_ttl_frames(&mut w, src, crate::AofFormat::V2, scratch)?;
crate::segmented::write_segmented_frames(&mut w, src, cold_seqs, scratch)?;
w.flush()?;
let inner = w.into_inner().map_err(|e| io::Error::other(e.to_string()))?;
let bytes = inner.metadata().map_or(0, |m| m.len());
inner.sync_all()?;
Ok((keys, bytes))
}
fn write_hash_ttl_frames<W: Write, S: crate::SnapshotSource>(
w: &mut W,
src: &S,
fmt: crate::AofFormat,
scratch: &mut Vec<u8>,
) -> io::Result<()> {
let mut err: Option<io::Error> = None;
src.for_each_hash_ttl(|key, field, deadline_ms| {
if err.is_some() {
return;
}
let ms = deadline_ms.to_string();
let mut argv = kevy_resp::Argv::with_capacity(6, 0);
argv.push(b"HPEXPIREAT");
argv.push(key);
argv.push(ms.as_bytes());
argv.push(b"FIELDS");
argv.push(b"1");
argv.push(field);
if let Err(e) = emit(w, &argv, fmt, scratch) {
err = Some(e);
}
});
match err {
Some(e) => Err(e),
None => Ok(()),
}
}
pub fn dump_store_to_buf<S: crate::SnapshotSource>(
src: &S,
fmt: crate::AofFormat,
) -> (Vec<u8>, u64) {
let mut buf = Vec::with_capacity(crate::SNAPSHOT_BUF_CAP);
buf.extend_from_slice(match fmt {
crate::AofFormat::V1 => crate::aof::AOF_MAGIC,
crate::AofFormat::V2 => crate::record::AOF2_MAGIC,
});
let mut scratch = Vec::new();
let mut keys = 0u64;
src.for_each_entry(|key, value, ttl_ms| {
let _ = write_value_as_commands(&mut buf, key, value, ttl_ms, fmt, &mut scratch);
keys += 1;
});
(buf, keys)
}
#[cold]
pub(crate) fn write_value_as_commands<W: Write>(
w: &mut W,
key: &[u8],
value: &Value,
ttl_ms: Option<u64>,
fmt: crate::AofFormat,
scratch: &mut Vec<u8>,
) -> io::Result<()> {
match value {
Value::Cold(_) => {
return Err(io::Error::other(
"cold stub reached the AOF rewrite — SnapshotSource must materialize (T4)",
));
}
Value::Str(s) => write_verb_items(w, b"SET", key, 1, [s.to_vec()], fmt, scratch)?,
Value::Int(n) => {
write_verb_items(w, b"SET", key, 1, [n.to_string().into_bytes()], fmt, scratch)?
}
Value::ArcBulk(a) => {
write_verb_items(w, b"SET", key, 1, [a.as_ref().to_vec()], fmt, scratch)?
}
Value::Hash(h) => {
let fv = h.iter().flat_map(|(f, v)| [f.to_vec(), v.to_vec()]);
write_verb_items(w, b"HSET", key, 2, fv, fmt, scratch)?;
}
Value::SmallHashInline(h) => {
let fv = h.iter().flat_map(|(f, v)| [f.to_vec(), v.to_vec()]);
write_verb_items(w, b"HSET", key, 2, fv, fmt, scratch)?;
}
Value::PackedRow(r) => {
let fv = r
.fields()
.map(|(f, v)| (f.to_vec(), v.to_vec()))
.collect::<Vec<_>>()
.into_iter()
.flat_map(|(f, v)| [f, v]);
write_verb_items(w, b"HSET", key, 2, fv, fmt, scratch)?;
}
Value::SegHash(h) => {
let fv = h.iter().flat_map(|(f, v)| [f.to_vec(), v.to_vec()]);
write_verb_items(w, b"HSET", key, 2, fv, fmt, scratch)?;
}
Value::List(l) => write_verb_items(w, b"RPUSH", key, 1, l.iter().cloned(), fmt, scratch)?,
Value::SegList(l) => {
write_verb_items(w, b"RPUSH", key, 1, l.iter().cloned(), fmt, scratch)?;
}
Value::SmallListInline(l) => {
write_verb_items(w, b"RPUSH", key, 1, l.iter().map(<[u8]>::to_vec), fmt, scratch)?;
}
Value::Set(s) => {
write_verb_items(
w,
b"SADD",
key,
1,
s.iter().map(kevy_store::SmallBytes::to_vec),
fmt,
scratch,
)?;
}
Value::SegSet(s) => {
write_verb_items(
w,
b"SADD",
key,
1,
s.keys().map(kevy_store::SmallBytes::to_vec),
fmt,
scratch,
)?;
}
Value::SmallSetInline(s) => {
write_verb_items(w, b"SADD", key, 1, s.iter().map(<[u8]>::to_vec), fmt, scratch)?;
}
Value::ZSet(z) => {
let ms = z.ordered().flat_map(|(m, sc)| [fmt_zset_score(sc), m.to_vec()]);
write_verb_items(w, b"ZADD", key, 2, ms, fmt, scratch)?;
}
Value::SegZSet(z) => {
let ms = z.ordered().flat_map(|(m, sc)| [fmt_zset_score(sc), m.to_vec()]);
write_verb_items(w, b"ZADD", key, 2, ms, fmt, scratch)?;
}
Value::SmallZSetInline(z) => {
let ms = z.iter().flat_map(|(m, sc)| [fmt_zset_score(sc), m.to_vec()]);
write_verb_items(w, b"ZADD", key, 2, ms, fmt, scratch)?;
}
Value::Stream(s) => {
_ = crate::rewrite_stream_fmt::stream_as_commands(w, key, s, fmt, scratch)?
}
}
write_pexpireat(w, key, ttl_ms, fmt, scratch)
}
fn write_pexpireat<W: Write>(
w: &mut W,
key: &[u8],
ttl_ms: Option<u64>,
fmt: crate::AofFormat,
scratch: &mut Vec<u8>,
) -> io::Result<()> {
let Some(ms) = ttl_ms else { return Ok(()) };
let deadline = kevy_store::now_unix_ms().saturating_add(ms);
write_verb_items(w, b"PEXPIREAT", key, 1, [deadline.to_string().into_bytes()], fmt, scratch)
}
fn fmt_zset_score(s: f64) -> Vec<u8> {
#[allow(clippy::float_cmp)]
let is_integer_valued = s.is_finite() && s == s.trunc();
if is_integer_valued && s.abs() < 1e17 {
format!("{}", s as i64).into_bytes()
} else {
format!("{s:.17}").into_bytes()
}
}
pub(crate) fn estimate_multibulk_bytes<A: ArgvView + ?Sized>(args: &A) -> u64 {
let mut n: u64 = 3 + u64::from(decimal_digits(args.len() as u64));
for i in 0..args.len() {
let a = &args[i];
n += 3 + u64::from(decimal_digits(a.len() as u64)) + a.len() as u64 + 2;
}
n
}
#[inline]
fn decimal_digits(mut x: u64) -> u32 {
if x == 0 {
return 1;
}
let mut d = 0;
while x > 0 {
d += 1;
x /= 10;
}
d
}
pub fn write_multibulk<W: Write, A: ArgvView + ?Sized>(w: &mut W, args: &A) -> io::Result<()> {
write!(w, "*{}\r\n", args.len())?;
for i in 0..args.len() {
let a = &args[i];
write!(w, "${}\r\n", a.len())?;
w.write_all(a)?;
w.write_all(b"\r\n")?;
}
Ok(())
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) const REWRITE_EMIT_VERBS: &[&str] = &[
"SET",
"HSET",
"RPUSH",
"SADD",
"ZADD",
"PEXPIREAT",
"HPEXPIREAT",
"XADD",
"XSETID",
"XGROUP",
"XCLAIM",
];
#[cfg(test)]
mod op_table_parity {
use super::REWRITE_EMIT_VERBS;
use kevy_resp::ops_table::{ops_with, surface};
use std::collections::BTreeSet;
#[test]
fn rewrite_manifest_matches_table() {
let m: BTreeSet<&str> = REWRITE_EMIT_VERBS.iter().copied().collect();
let t: BTreeSet<&str> = ops_with(surface::REWRITE).into_iter().collect();
assert_eq!(
m, t,
"rewrite emit set != OP_TABLE REWRITE flags — update both when changing rewrite_fmt"
);
}
#[test]
fn rewrite_manifest_verbs_have_source_literals() {
let src = include_str!("rewrite_fmt.rs");
for v in REWRITE_EMIT_VERBS {
let lit = format!("\"{v}");
assert!(
src.contains(&lit) || src.contains(&format!("b\"{v}\"")),
"REWRITE_EMIT_VERBS lists {v} but rewrite_fmt.rs has no literal for it"
);
}
}
}