use yo_common::crc::crc64;
use crate::db::Db;
use crate::keys::Record;
use crate::rdb;
use crate::value::Kind;
use yo_index::Cursor as KeyCursor;
const HEADER: &[u8] = b"REDIS0012";
const OP_AUX: u8 = 0xFA;
const OP_RESIZEDB: u8 = 0xFB;
const OP_EXPIRETIME_MS: u8 = 0xFC;
const OP_SELECTDB: u8 = 0xFE;
const OP_EOF: u8 = 0xFF;
const BATCH: usize = 256;
pub struct Snapshot {
out: Vec<u8>,
crc: u64,
taken: usize,
started: bool,
names: Vec<u8>,
bounds: Vec<(usize, usize)>,
skipped: usize,
}
impl Default for Snapshot {
fn default() -> Snapshot {
Snapshot::new()
}
}
impl Snapshot {
#[must_use]
pub fn new() -> Snapshot {
let mut snap = Snapshot {
out: Vec::new(),
crc: 0,
taken: 0,
started: false,
names: Vec::new(),
bounds: Vec::new(),
skipped: 0,
};
snap.out.extend_from_slice(HEADER);
snap.aux(b"yo-ver", env!("CARGO_PKG_VERSION").as_bytes());
snap
}
pub fn aux(&mut self, name: &[u8], value: &[u8]) {
debug_assert!(!self.started, "an aux field after the first database");
self.out.push(OP_AUX);
rdb::put_str(&mut self.out, name);
rdb::put_str(&mut self.out, value);
self.absorb();
}
pub fn database(&mut self, index: usize, db: &mut Db) {
if db.is_empty() {
return;
}
self.started = true;
self.out.push(OP_SELECTDB);
rdb::put_len(&mut self.out, index as u64);
self.out.push(OP_RESIZEDB);
rdb::put_len(&mut self.out, db.len() as u64);
rdb::put_len(&mut self.out, db.expires() as u64);
self.absorb();
let mut at = KeyCursor::START;
loop {
self.names.clear();
self.bounds.clear();
let names = &mut self.names;
let bounds = &mut self.bounds;
at = db.scan(at, BATCH, None, |key| {
let from = names.len();
names.extend_from_slice(key);
bounds.push((from, names.len()));
});
for i in 0..self.bounds.len() {
let (from, to) = self.bounds[i];
let key = self.names[from..to].to_vec();
let Some(rec) = db.at(&key).export(&key) else {
continue;
};
self.entry(&key, &rec);
}
if at.is_end() {
break;
}
}
}
#[must_use]
pub fn finish(mut self) -> Vec<u8> {
self.out.push(OP_EOF);
self.absorb();
self.out.extend_from_slice(&self.crc.to_le_bytes());
self.out
}
#[must_use]
pub const fn skipped(&self) -> usize {
self.skipped
}
fn entry(&mut self, key: &[u8], rec: &Record) {
if matches!(rec.kind(), Kind::Foreign | Kind::Array) {
self.skipped += 1;
return;
}
if let Some(at) = rec.expire_at() {
self.out.push(OP_EXPIRETIME_MS);
self.out.extend_from_slice(&at.to_le_bytes());
}
let wrote = rdb::object(rec, Some(key), &mut self.out);
debug_assert!(wrote, "the kinds without a shape were refused above");
self.absorb();
}
fn absorb(&mut self) {
self.crc = crc64(self.crc, &self.out[self.taken..]);
self.taken = self.out.len();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::clock::Clock;
use crate::lists::End;
use crate::rdb::FOOTER;
use crate::streams::{Add, Trim};
use crate::ttl::Cond;
use crate::zsets::ZAdd;
struct Parse<'a> {
buf: &'a [u8],
at: usize,
}
impl<'a> Parse<'a> {
fn byte(&mut self) -> u8 {
let b = self.buf[self.at];
self.at += 1;
b
}
fn len(&mut self) -> u64 {
let first = self.byte();
match first >> 6 {
0 => u64::from(first & 0x3f),
1 => (u64::from(first & 0x3f) << 8) | u64::from(self.byte()),
_ => panic!("the tests do not write a length that big"),
}
}
fn str(&mut self) -> Vec<u8> {
let n = self.len() as usize;
let s = self.buf[self.at..self.at + n].to_vec();
self.at += n;
s
}
}
#[derive(Default, PartialEq, Eq, Debug)]
struct Seen {
aux: Vec<(Vec<u8>, Vec<u8>)>,
dbs: Vec<u64>,
keys: Vec<(Vec<u8>, Option<u64>, Vec<u8>)>,
}
fn walk(file: &[u8]) -> Seen {
assert_eq!(&file[..9], HEADER, "the header");
let crc = u64::from_le_bytes(file[file.len() - 8..].try_into().unwrap());
assert_eq!(crc, crc64(0, &file[..file.len() - 8]), "the checksum");
let mut seen = Seen::default();
let mut p = Parse {
buf: file,
at: HEADER.len(),
};
let mut expire = None;
loop {
match p.byte() {
OP_AUX => {
let name = p.str();
let value = p.str();
seen.aux.push((name, value));
}
OP_SELECTDB => seen.dbs.push(p.len()),
OP_RESIZEDB => {
p.len();
p.len();
}
OP_EXPIRETIME_MS => {
let at = u64::from_le_bytes(p.buf[p.at..p.at + 8].try_into().unwrap());
p.at += 8;
expire = Some(at);
}
OP_EOF => {
assert_eq!(p.at + 8, file.len(), "the end is where the file ends");
return seen;
}
ty => {
let key = p.str();
let mut payload = vec![ty];
let rest = &p.buf[p.at..];
let taken = crate::rdb::measure(ty, rest).expect("a value we wrote");
payload.extend_from_slice(&rest[..taken]);
p.at += taken;
seen.keys.push((key, expire.take(), payload));
}
}
}
}
fn db() -> Db {
Db::with_clock(Clock::fixed(1_000), 1)
}
#[test]
fn an_empty_server_is_a_header_an_aux_field_and_an_end() {
let mut snap = Snapshot::new();
snap.database(0, &mut db());
let file = snap.finish();
let seen = walk(&file);
assert!(seen.dbs.is_empty(), "no database has anything in it");
assert!(seen.keys.is_empty());
assert_eq!(seen.aux.len(), 1, "the one we write");
}
#[test]
fn every_type_goes_out_as_the_payload_dump_would_have_written() {
let mut d = db();
d.at(b"str").set_plain(b"str", b"hello").unwrap();
d.at(b"list")
.push(b"list", End::Right, [b"a".as_slice(), b"b"].into_iter())
.unwrap();
d.at(b"set")
.sadd(b"set", [b"x".as_slice(), b"y"].into_iter())
.unwrap();
d.at(b"ints")
.sadd(b"ints", [b"1".as_slice(), b"2"].into_iter())
.unwrap();
d.at(b"zset")
.zadd(
b"zset",
[(1.5, b"m".as_slice())].into_iter(),
ZAdd::default(),
)
.unwrap();
d.at(b"hash")
.hset(b"hash", [(b"f".as_slice(), b"v".as_slice())].into_iter())
.unwrap();
d.at(b"stream")
.xadd(
b"stream",
Add::Auto,
&[(b"f".as_slice(), b"v".as_slice())],
Trim::None,
true,
1_000,
)
.unwrap();
let mut snap = Snapshot::new();
snap.database(0, &mut d);
let file = snap.finish();
let seen = walk(&file);
assert_eq!(seen.dbs, vec![0]);
assert_eq!(seen.keys.len(), 7, "one entry per key");
for (key, expire, payload) in &seen.keys {
assert_eq!(*expire, None, "nothing was given a deadline");
assert_eq!(
Some(payload.clone()),
d.at(key).dump(key).map(|p| p[..p.len() - FOOTER].to_vec()),
"{}",
String::from_utf8_lossy(key)
);
}
}
#[test]
fn a_deadline_travels_with_the_key_and_a_dead_key_does_not() {
let mut d = db();
for key in [&b"alive"[..], b"soon", b"gone"] {
d.at(key).set_plain(key, b"1").unwrap();
}
d.at(b"soon").expire(b"soon", 9_000, Cond::Always);
d.at(b"gone").expire(b"gone", 500, Cond::Always);
let mut snap = Snapshot::new();
snap.database(0, &mut d);
let seen = walk(&snap.finish());
let mut found: Vec<(Vec<u8>, Option<u64>)> = seen
.keys
.iter()
.map(|(k, at, _)| (k.clone(), *at))
.collect();
found.sort();
assert_eq!(
found,
vec![(b"alive".to_vec(), None), (b"soon".to_vec(), Some(9_000))],
"the key whose deadline has gone is not in the file"
);
}
#[test]
fn each_database_is_selected_once_and_the_empty_ones_are_not_there() {
let mut zero = db();
let mut empty = db();
let mut nine = db();
zero.at(b"a").set_plain(b"a", b"1").unwrap();
nine.at(b"b").set_plain(b"b", b"2").unwrap();
let mut snap = Snapshot::new();
snap.database(0, &mut zero);
snap.database(4, &mut empty);
snap.database(9, &mut nine);
let seen = walk(&snap.finish());
assert_eq!(seen.dbs, vec![0, 9], "the empty one in the middle is gone");
assert_eq!(seen.keys.len(), 2);
}
#[test]
fn more_keys_than_one_batch_all_arrive_once() {
let mut d = db();
let many = BATCH * 3 + 7;
for i in 0..many {
d.at(format!("k{i}").as_bytes())
.set_plain(format!("k{i}").as_bytes(), b"v")
.unwrap();
}
let mut snap = Snapshot::new();
snap.database(0, &mut d);
let seen = walk(&snap.finish());
let mut names: Vec<&Vec<u8>> = seen.keys.iter().map(|(k, _, _)| k).collect();
names.sort();
names.dedup();
assert_eq!(names.len(), many, "every key, and none of them twice");
}
#[test]
fn a_striped_database_writes_every_key_once_and_one_selector() {
let mut d = Db::with_clock(Clock::fixed(1_000), 8);
let many = BATCH * 3 + 7;
for i in 0..many {
let key = format!("k{i}");
d.at(key.as_bytes())
.set_plain(key.as_bytes(), b"v")
.unwrap();
}
assert!(
d.stripes_mut().all(|s| !s.is_empty()),
"a stripe got nothing, so this proves less than it looks"
);
d.at(b"k5").expire(b"k5", 9_000, Cond::Always);
let mut snap = Snapshot::new();
snap.database(3, &mut d);
let seen = walk(&snap.finish());
assert_eq!(seen.dbs, vec![3], "one selector for the whole database");
let mut names: Vec<&Vec<u8>> = seen.keys.iter().map(|(k, _, _)| k).collect();
names.sort();
names.dedup();
assert_eq!(names.len(), many, "every key, and none of them twice");
let dated: Vec<&Vec<u8>> = seen
.keys
.iter()
.filter(|(_, e, _)| e.is_some())
.map(|(k, _, _)| k)
.collect();
assert_eq!(
dated,
vec![&b"k5".to_vec()],
"the deadline went on that key"
);
}
#[test]
fn a_type_with_no_rdb_shape_is_counted_and_left_out() {
let mut d = db();
d.at(b"ordinary").set_plain(b"ordinary", b"1").unwrap();
d.at(b"sparse")
.arset(b"sparse", 7, [&b"a"[..]].into_iter())
.unwrap();
let mut snap = Snapshot::new();
snap.database(0, &mut d);
assert_eq!(snap.skipped(), 1);
let seen = walk(&snap.finish());
assert_eq!(seen.keys.len(), 1);
assert_eq!(seen.keys[0].0, b"ordinary");
}
#[test]
fn an_aux_field_the_caller_adds_is_in_the_file() {
let mut snap = Snapshot::new();
snap.aux(b"redis-ver", b"8.8.0");
let seen = walk(&snap.finish());
assert!(
seen.aux
.contains(&(b"redis-ver".to_vec(), b"8.8.0".to_vec())),
"{:?}",
seen.aux
);
}
}