use yo_common::Result;
use crate::array::Array;
use crate::hash::Hash;
use crate::keyspace::Keyspace;
use crate::list::List;
use crate::rdb;
use crate::set::Set;
use crate::value::{self, Kind};
use crate::zset::Zset;
#[derive(Debug, Clone)]
pub struct Record {
body: Body,
expire_at: Option<u64>,
}
impl Record {
pub(crate) const fn new(body: Body, expire_at: Option<u64>) -> Record {
Record { body, expire_at }
}
pub(crate) const fn body(&self) -> &Body {
&self.body
}
#[must_use]
pub const fn kind(&self) -> Kind {
match self.body {
Body::String(_) => Kind::String,
Body::Set(_) => Kind::Set,
Body::Hash(_) => Kind::Hash,
Body::List(_) => Kind::List,
Body::Zset(_) => Kind::Zset,
Body::Array(_) => Kind::Array,
}
}
#[must_use]
pub const fn expire_at(&self) -> Option<u64> {
self.expire_at
}
}
#[derive(Debug, Clone)]
pub(crate) enum Body {
String(Vec<u8>),
Set(Set),
Hash(Hash),
List(List),
Zset(Zset),
Array(Array),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Moved {
Missing,
Taken,
Ok,
}
impl Keyspace {
pub fn export(&mut self, key: &[u8]) -> Option<Record> {
let addr = self.live_rec(key)?;
let rec = self.map.value_at(addr);
let expire_at = value::expire_at(rec);
let body = match value::kind(rec) {
Kind::String => Body::String(value::read(rec).to_vec()),
Kind::Set => Body::Set(
self.sets
.get(value::slot(rec))
.expect("the record points at its body")
.clone(),
),
Kind::Hash => Body::Hash(
self.hashes
.get(value::slot(rec))
.expect("the record points at its body")
.clone(),
),
Kind::List => Body::List(
self.lists
.get(value::slot(rec))
.expect("the record points at its body")
.clone(),
),
Kind::Zset => Body::Zset(
self.zsets
.get(value::slot(rec))
.expect("the record points at its body")
.clone(),
),
Kind::Array => Body::Array(
self.arrays
.get(value::slot(rec))
.expect("the record points at its body")
.clone(),
),
Kind::Stream => unreachable!("nothing can store a stream yet"),
};
Some(Record { body, expire_at })
}
pub fn take(&mut self, key: &[u8]) -> Option<Record> {
let addr = self.live_rec(key)?;
let rec = self.map.value_at(addr);
let expire_at = value::expire_at(rec);
let kind = value::kind(rec);
if kind == Kind::String {
let bytes = value::read(rec).to_vec();
self.del_rec(key);
return Some(Record {
body: Body::String(bytes),
expire_at,
});
}
let slot = value::slot(rec);
let gone = "the record points at its body";
let body = match kind {
Kind::Set => Body::Set(self.sets.remove(slot).expect(gone)),
Kind::Hash => Body::Hash(self.hashes.remove(slot).expect(gone)),
Kind::List => Body::List(self.lists.remove(slot).expect(gone)),
Kind::Zset => Body::Zset(self.zsets.remove(slot).expect(gone)),
Kind::Array => Body::Array(self.arrays.remove(slot).expect(gone)),
Kind::String | Kind::Stream => unreachable!("handled above or cannot be stored"),
};
self.bodies -= 1;
self.del_rec(key);
Some(Record { body, expire_at })
}
pub fn import(&mut self, key: &[u8], rec: Record) {
let at = rec.expire_at;
match rec.body {
Body::String(bytes) => self.store(key, &bytes, at),
Body::Set(set) => {
self.free_body(key);
let slot = self.sets.insert(set);
self.bodies += 1;
self.write_slot(key, Kind::Set, slot, at);
}
Body::Hash(hash) => {
self.free_body(key);
let slot = self.hashes.insert(hash);
self.bodies += 1;
self.write_slot(key, Kind::Hash, slot, at);
}
Body::List(list) => {
self.free_body(key);
let slot = self.lists.insert(list);
self.bodies += 1;
self.write_slot(key, Kind::List, slot, at);
}
Body::Zset(zset) => {
self.free_body(key);
let slot = self.zsets.insert(zset);
self.bodies += 1;
self.write_slot(key, Kind::Zset, slot, at);
}
Body::Array(array) => {
self.free_body(key);
let slot = self.arrays.insert(array);
self.bodies += 1;
self.write_slot(key, Kind::Array, slot, at);
}
}
}
pub fn dump(&mut self, key: &[u8]) -> Option<Vec<u8>> {
let rec = self.export(key)?;
rdb::dump(&rec)
}
pub fn restore(
&mut self,
key: &[u8],
payload: &[u8],
expire_at: Option<u64>,
replace: bool,
) -> std::result::Result<Moved, rdb::Bad> {
if !replace && self.exists(key) {
return Ok(Moved::Taken);
}
let limits = rdb::Limits {
set: &self.limits,
hash: &self.hash_limits,
list: &self.list_limits,
zset: &self.zset_limits,
};
let now = self.clock.now_ms();
let body = rdb::load(payload, limits, now)?;
if expire_at.is_some_and(|at| at <= now) {
self.del(key);
return Ok(Moved::Ok);
}
self.import(key, Record::new(body, expire_at));
Ok(Moved::Ok)
}
pub fn rename(&mut self, src: &[u8], dst: &[u8], only_if_new: bool) -> Moved {
if self.live_rec(src).is_none() {
return Moved::Missing;
}
let same = src == dst;
if only_if_new && (same || self.live_rec(dst).is_some()) {
return Moved::Taken;
}
if same {
return Moved::Ok;
}
let addr = self.map.find(src).expect("it was live a line ago");
let mut bytes = std::mem::take(&mut self.scratch);
bytes.clear();
bytes.extend_from_slice(self.map.value_at(addr));
self.free_body(dst);
self.write_rec(dst, bytes.len(), |out| {
out.copy_from_slice(&bytes);
});
self.scratch = bytes;
self.del_rec(src);
Moved::Ok
}
pub fn copy(&mut self, src: &[u8], dst: &[u8], replace: bool) -> Moved {
if self.live_rec(src).is_none() {
return Moved::Missing;
}
let same = src == dst;
if !replace && (same || self.live_rec(dst).is_some()) {
return Moved::Taken;
}
if same {
return Moved::Ok;
}
let addr = self.map.find(src).expect("it was live a line ago");
if value::kind(self.map.value_at(addr)) == Kind::String {
let mut bytes = std::mem::take(&mut self.scratch);
bytes.clear();
bytes.extend_from_slice(self.map.value_at(addr));
self.free_body(dst);
self.write_rec(dst, bytes.len(), |out| {
out.copy_from_slice(&bytes);
});
self.scratch = bytes;
return Moved::Ok;
}
let rec = self.export(src).expect("it was live a line ago");
self.import(dst, rec);
Moved::Ok
}
pub fn touch<'k>(&mut self, keys: impl Iterator<Item = &'k [u8]>) -> usize {
keys.filter(|key| self.exists(key)).count()
}
fn write_slot(&mut self, key: &[u8], kind: Kind, slot: u32, at: Option<u64>) {
let len = value::slot_record_len(at.is_some());
self.write_rec(key, len, |out| {
value::write_slot_record(out, kind, slot, at);
});
}
}
#[must_use]
pub fn no_such_key() -> yo_common::Error {
yo_common::Error::new(yo_common::Code::Invalid, "no such key")
}
impl Moved {
pub fn found(self) -> Result<Moved> {
match self {
Moved::Missing => Err(no_such_key()),
other => Ok(other),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Clock;
use crate::End;
use crate::zsets::ZAdd;
use crate::{Applied, Cond};
fn db() -> Keyspace {
Keyspace::with_clock(Clock::fixed(1_000_000))
}
fn members(d: &mut Keyspace, key: &[u8]) -> Vec<String> {
let mut out: Vec<String> = d
.smembers(key)
.expect("a set")
.expect("a key")
.map(|m| String::from_utf8(m.to_vec()).expect("utf8 in these tests"))
.collect();
out.sort();
out
}
fn put(d: &mut Keyspace, key: &[u8], val: &[u8]) {
d.set_plain(key, val).expect("room for a record");
}
fn read(d: &mut Keyspace, key: &[u8]) -> Vec<u8> {
d.get(key).expect("a string").expect("there").to_vec()
}
#[test]
fn a_rename_moves_the_value_and_leaves_nothing_behind() {
let mut d = db();
put(&mut d, b"a", b"v1");
assert_eq!(d.rename(b"a", b"b", false), Moved::Ok);
assert!(!d.exists(b"a"));
assert_eq!(read(&mut d, b"b"), b"v1");
}
#[test]
fn a_rename_does_not_allocate_to_carry_the_record_across() {
let mut d = db();
put(&mut d, b"a", b"v1");
for _ in 0..4 {
assert_eq!(d.rename(b"a", b"b", false), Moved::Ok);
assert_eq!(d.rename(b"b", b"a", false), Moved::Ok);
}
let (_, allocs) = crate::tally::counted(|| {
for _ in 0..50 {
assert_eq!(d.rename(b"a", b"b", false), Moved::Ok);
assert_eq!(d.rename(b"b", b"a", false), Moved::Ok);
}
});
assert_eq!(allocs, 0, "rename allocated {allocs} times in a hundred");
assert_eq!(read(&mut d, b"a"), b"v1");
}
#[test]
fn a_rename_with_no_source_is_the_one_error_in_this_file() {
let mut d = db();
assert_eq!(d.rename(b"a", b"b", false), Moved::Missing);
assert_eq!(d.rename(b"a", b"b", true), Moved::Missing);
assert_eq!(
d.copy(b"a", b"b", false),
Moved::Missing,
"copy just says 0"
);
}
#[test]
fn a_rename_carries_the_source_deadline_and_drops_the_destination_one() {
let mut d = db();
put(&mut d, b"a", b"v1");
d.set_expiry(b"a", Some(2_000_000));
put(&mut d, b"b", b"v2");
d.set_expiry(b"b", Some(1_500_000));
assert_eq!(d.rename(b"a", b"b", false), Moved::Ok);
assert_eq!(d.deadline_of(b"b"), crate::Ask::At(2_000_000));
}
#[test]
fn renaming_a_key_onto_itself_keeps_it_and_renamenx_refuses() {
let mut d = db();
put(&mut d, b"a", b"v1");
d.set_expiry(b"a", Some(2_000_000));
assert_eq!(d.rename(b"a", b"a", false), Moved::Ok);
assert_eq!(read(&mut d, b"a"), b"v1");
assert_eq!(d.deadline_of(b"a"), crate::Ask::At(2_000_000));
assert_eq!(d.rename(b"a", b"a", true), Moved::Taken);
}
#[test]
fn renamenx_writes_over_nothing() {
let mut d = db();
put(&mut d, b"a", b"v1");
put(&mut d, b"b", b"v2");
assert_eq!(d.rename(b"a", b"b", true), Moved::Taken);
assert_eq!(read(&mut d, b"a"), b"v1");
assert_eq!(read(&mut d, b"b"), b"v2");
assert_eq!(d.rename(b"a", b"c", true), Moved::Ok);
assert!(!d.exists(b"a"));
}
#[test]
fn renaming_a_set_moves_the_slot_and_not_the_members() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref(), b"m2".as_ref()].into_iter())
.expect("a set");
let before = d.memory_bytes();
assert_eq!(d.rename(b"s", b"t", false), Moved::Ok);
assert_eq!(members(&mut d, b"t"), ["m1", "m2"]);
assert_eq!(d.kind_of(b"t"), Some(Kind::Set));
assert!(!d.exists(b"s"));
assert!(
d.memory_bytes().abs_diff(before) < 64,
"the members were not copied"
);
}
#[test]
fn renaming_over_a_set_frees_the_set_that_was_there() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref()].into_iter()).expect("a set");
d.sadd(b"t", [b"m2".as_ref()].into_iter()).expect("a set");
assert_eq!(d.sets.len(), 2);
assert_eq!(d.rename(b"s", b"t", false), Moved::Ok);
assert_eq!(d.sets.len(), 1, "the destination's body went with it");
assert_eq!(members(&mut d, b"t"), ["m1"]);
}
#[test]
fn a_copy_is_a_second_value_and_not_a_second_name() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref(), b"m2".as_ref()].into_iter())
.expect("a set");
assert_eq!(d.copy(b"s", b"t", false), Moved::Ok);
d.sadd(b"t", [b"m3".as_ref()].into_iter()).expect("a set");
assert_eq!(
members(&mut d, b"s"),
["m1", "m2"],
"the original is intact"
);
assert_eq!(members(&mut d, b"t"), ["m1", "m2", "m3"]);
}
#[test]
fn a_copy_refuses_a_destination_it_was_not_told_it_could_have() {
let mut d = db();
put(&mut d, b"a", b"v1");
put(&mut d, b"b", b"v2");
assert_eq!(d.copy(b"a", b"b", false), Moved::Taken);
assert_eq!(read(&mut d, b"b"), b"v2");
assert_eq!(d.copy(b"a", b"b", true), Moved::Ok);
assert_eq!(read(&mut d, b"b"), b"v1");
}
#[test]
fn a_copy_of_a_string_does_not_allocate() {
let mut d = db();
put(&mut d, b"a", b"a-value-of-some-length");
for _ in 0..4 {
assert_eq!(d.copy(b"a", b"b", true), Moved::Ok);
}
let (_, allocs) = crate::tally::counted(|| {
for _ in 0..50 {
assert_eq!(d.copy(b"a", b"b", true), Moved::Ok);
}
});
assert_eq!(allocs, 0, "copy allocated {allocs} times in fifty");
assert_eq!(read(&mut d, b"b"), b"a-value-of-some-length");
}
#[test]
fn a_copy_onto_itself_leaves_the_key_alone() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref(), b"m2".as_ref()].into_iter())
.expect("a set");
assert_eq!(d.copy(b"s", b"s", false), Moved::Taken);
assert_eq!(d.copy(b"s", b"s", true), Moved::Ok);
assert_eq!(members(&mut d, b"s"), ["m1", "m2"]);
assert_eq!(d.sets.len(), 1, "no second body was made or lost");
}
#[test]
fn a_copy_carries_the_deadline() {
let mut d = db();
put(&mut d, b"a", b"v1");
d.set_expiry(b"a", Some(2_000_000));
assert_eq!(d.copy(b"a", b"b", false), Moved::Ok);
assert_eq!(d.deadline_of(b"b"), crate::Ask::At(2_000_000));
assert_eq!(d.deadline_of(b"a"), crate::Ask::At(2_000_000));
}
#[test]
fn a_destination_that_has_already_gone_counts_as_free() {
let mut d = db();
put(&mut d, b"a", b"v1");
put(&mut d, b"b", b"v2");
d.set_expiry(b"b", Some(999_999));
assert_eq!(d.copy(b"a", b"b", false), Moved::Ok, "b was already gone");
assert_eq!(read(&mut d, b"b"), b"v1");
}
#[test]
fn a_source_that_has_already_gone_is_not_a_source() {
let mut d = db();
put(&mut d, b"a", b"v1");
d.set_expiry(b"a", Some(999_999));
assert_eq!(d.rename(b"a", b"b", false), Moved::Missing);
assert_eq!(d.copy(b"a", b"b", false), Moved::Missing);
}
#[test]
fn a_record_taken_out_of_a_database_outlives_it() {
let mut from = db();
from.sadd(b"s", [b"m1".as_ref(), b"m2".as_ref()].into_iter())
.expect("a set");
let rec = from.export(b"s").expect("a record");
assert_eq!(rec.kind(), Kind::Set);
from.clear();
let mut into = db();
into.import(b"s", rec);
assert_eq!(members(&mut into, b"s"), ["m1", "m2"]);
}
#[test]
fn importing_over_a_body_does_not_leave_it_in_the_slab() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref()].into_iter()).expect("a set");
d.sadd(b"t", [b"m2".as_ref()].into_iter()).expect("a set");
let rec = d.export(b"s").expect("a record");
d.import(b"t", rec);
assert_eq!(d.sets.len(), 2, "s and t, and not the one t used to hold");
assert_eq!(members(&mut d, b"t"), ["m1"]);
}
#[test]
fn importing_a_string_over_a_set_frees_the_set() {
let mut d = db();
put(&mut d, b"a", b"v1");
d.sadd(b"s", [b"m1".as_ref()].into_iter()).expect("a set");
assert_eq!(d.sets.len(), 1);
assert_eq!(d.copy(b"a", b"s", true), Moved::Ok);
assert_eq!(d.sets.len(), 0, "the set went when the string arrived");
assert_eq!(d.kind_of(b"s"), Some(Kind::String));
}
#[test]
fn a_list_can_be_copied_and_the_copy_is_its_own() {
let mut d = db();
d.push(b"l", End::Left, [b"a".as_ref(), b"b".as_ref()].into_iter())
.expect("a list");
assert_eq!(d.copy(b"l", b"m", false), Moved::Ok);
assert_eq!(d.kind_of(b"m"), Some(Kind::List));
assert_eq!(d.llen(b"m").expect("a list"), 2);
d.push(b"m", End::Left, [b"c".as_ref()].into_iter())
.expect("a list");
assert_eq!(d.llen(b"l").expect("a list"), 2, "the source did not grow");
assert_eq!(d.llen(b"m").expect("a list"), 3);
}
#[test]
fn a_zset_can_be_copied_and_the_copy_is_its_own() {
let mut d = db();
d.zadd(b"z", [(1.0, b"m1".as_ref())].into_iter(), ZAdd::default())
.expect("a zset");
assert_eq!(d.copy(b"z", b"y", false), Moved::Ok);
assert_eq!(d.kind_of(b"y"), Some(Kind::Zset));
assert_eq!(d.zscore(b"y", b"m1").expect("a zset"), Some(1.0));
d.zadd(b"y", [(2.0, b"m2".as_ref())].into_iter(), ZAdd::default())
.expect("a zset");
assert_eq!(d.zcard(b"z").expect("a zset"), 1, "the source did not grow");
assert_eq!(d.zcard(b"y").expect("a zset"), 2);
}
#[test]
fn copying_over_a_list_frees_the_list() {
let mut d = db();
put(&mut d, b"a", b"v1");
d.push(b"l", End::Left, [b"x".as_ref()].into_iter())
.expect("a list");
assert_eq!(d.copy(b"a", b"l", true), Moved::Ok);
assert_eq!(d.kind_of(b"l"), Some(Kind::String));
assert_eq!(read(&mut d, b"l"), b"v1");
}
#[test]
fn taking_a_set_empties_the_slab_and_the_key() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref(), b"m2".as_ref()].into_iter())
.expect("a set");
assert_eq!(d.sets.len(), 1);
let rec = d.take(b"s").expect("a record");
assert_eq!(rec.kind(), Kind::Set);
assert_eq!(d.sets.len(), 0, "the body left with the record");
assert!(!d.exists(b"s"), "and so did the key");
let mut into = db();
into.import(b"s", rec);
assert_eq!(members(&mut into, b"s"), ["m1", "m2"]);
}
#[test]
fn taking_a_string_takes_the_record_with_it() {
let mut d = db();
put(&mut d, b"a", b"v1");
let rec = d.take(b"a").expect("a record");
assert_eq!(rec.kind(), Kind::String);
assert!(d.take(b"a").is_none());
assert_eq!(d.len(), 0);
}
#[test]
fn a_taken_key_keeps_the_time_it_had_left() {
let mut d = db();
put(&mut d, b"a", b"v1");
assert_eq!(d.expire(b"a", 2_000_000, Cond::Always), Applied::Ok);
let rec = d.take(b"a").expect("a record");
assert_eq!(rec.expire_at(), Some(2_000_000));
}
#[test]
fn a_dead_key_cannot_be_taken() {
let mut d = db();
d.sadd(b"s", [b"m1".as_ref()].into_iter()).expect("a set");
assert_eq!(d.expire(b"s", 1_000_001, Cond::Always), Applied::Ok);
d.clock_mut().advance(10);
assert!(d.take(b"s").is_none());
assert_eq!(d.sets.len(), 0, "and the body did not stay behind");
}
#[test]
fn touch_counts_the_way_exists_counts() {
let mut d = db();
put(&mut d, b"a", b"v1");
put(&mut d, b"b", b"v2");
assert_eq!(d.touch([b"a".as_ref()].into_iter()), 1);
assert_eq!(d.touch([b"a".as_ref(), b"b".as_ref()].into_iter()), 2);
assert_eq!(d.touch([b"a".as_ref(), b"a".as_ref()].into_iter()), 2);
assert_eq!(d.touch([b"a".as_ref(), b"z".as_ref()].into_iter()), 1);
assert_eq!(d.touch([b"z".as_ref()].into_iter()), 0);
}
}