use std::sync::atomic::{AtomicU64, Ordering};
use yo_common::{Addr, Code, Error, Result, Rng, bytes_eq};
use yo_index::RawMap;
use crate::Clock;
use crate::array::Array;
use crate::hash::{self, Hash};
use crate::list::{self, List};
use crate::set::{self, Set};
use crate::slab::Slab;
use crate::ttl::{self, Applied, Ask, Cond};
use crate::value::{self, Kind};
use crate::zset::{self, Zset};
pub struct Keyspace {
pub(crate) map: RawMap,
pub(crate) clock: Clock,
pub(crate) expired: u64,
pub(crate) sets: Slab<Set>,
pub(crate) hashes: Slab<Hash>,
pub(crate) lists: Slab<List>,
pub(crate) zsets: Slab<Zset>,
pub(crate) arrays: Slab<Array>,
pub(crate) bodies: usize,
pub(crate) limits: set::Limits,
pub(crate) hash_limits: hash::Limits,
pub(crate) list_limits: list::Limits,
pub(crate) zset_limits: zset::Limits,
pub(crate) rng: Rng,
memo: Memo,
pub(crate) scratch: Vec<u8>,
pub(crate) rows: Vec<usize>,
pub(crate) setops: crate::setops::Scratch,
}
const SCRATCH: usize = 1024;
struct Memo {
writes: u64,
live: bool,
kind: Kind,
slot: u32,
len: u8,
key: [u8; Memo::MAX],
}
impl Memo {
const MAX: usize = 32;
const fn empty() -> Memo {
Memo {
writes: 0,
live: false,
kind: Kind::String,
slot: 0,
len: 0,
key: [0; Memo::MAX],
}
}
#[inline]
fn get(&self, writes: u64, key: &[u8]) -> Option<(Kind, u32)> {
if !self.live || self.writes != writes || key.len() != self.len as usize {
return None;
}
bytes_eq(&self.key[..key.len()], key).then_some((self.kind, self.slot))
}
#[inline]
fn put(&mut self, writes: u64, key: &[u8], kind: Kind, slot: u32) {
if key.len() > Memo::MAX {
self.live = false;
return;
}
self.writes = writes;
self.live = true;
self.kind = kind;
self.slot = slot;
self.len = key.len() as u8;
self.key[..key.len()].copy_from_slice(key);
}
}
static MADE: AtomicU64 = AtomicU64::new(0);
impl Keyspace {
#[must_use]
pub fn new() -> Keyspace {
Keyspace::with_clock(Clock::system())
}
#[must_use]
pub fn with_clock(clock: Clock) -> Keyspace {
let made = MADE.fetch_add(1, Ordering::Relaxed);
Keyspace {
map: RawMap::new(),
clock,
expired: 0,
sets: Slab::new(),
hashes: Slab::new(),
lists: Slab::new(),
zsets: Slab::new(),
arrays: Slab::new(),
bodies: 0,
limits: set::Limits::DEFAULT,
hash_limits: hash::Limits::DEFAULT,
list_limits: list::Limits::default(),
zset_limits: zset::Limits::DEFAULT,
rng: Rng::new(clock.now_ms() ^ made.wrapping_mul(0x9e37_79b9_7f4a_7c15)),
memo: Memo::empty(),
scratch: Vec::with_capacity(SCRATCH),
rows: Vec::new(),
setops: crate::setops::Scratch::new(),
}
}
#[inline]
pub const fn seed(&mut self, seed: u64) {
self.rng = Rng::new(seed);
}
#[inline]
pub const fn limits(&self) -> &set::Limits {
&self.limits
}
#[inline]
pub const fn set_limits(&mut self, limits: set::Limits) {
self.limits = limits;
}
#[inline]
pub const fn hash_limits(&self) -> &hash::Limits {
&self.hash_limits
}
#[inline]
pub const fn set_hash_limits(&mut self, limits: hash::Limits) {
self.hash_limits = limits;
}
#[inline]
pub const fn list_limits(&self) -> &list::Limits {
&self.list_limits
}
#[inline]
pub const fn set_list_limits(&mut self, limits: list::Limits) {
self.list_limits = limits;
}
#[inline]
pub const fn zset_limits(&self) -> &zset::Limits {
&self.zset_limits
}
#[inline]
pub const fn set_zset_limits(&mut self, limits: zset::Limits) {
self.zset_limits = limits;
}
#[inline]
pub const fn clock(&self) -> &Clock {
&self.clock
}
#[inline]
pub const fn clock_mut(&mut self) -> &mut Clock {
&mut self.clock
}
#[inline]
pub const fn map(&self) -> &RawMap {
&self.map
}
#[inline]
pub fn len(&self) -> usize {
self.map.len()
}
#[inline]
pub fn is_empty(&self) -> bool {
self.map.is_empty()
}
pub fn kind_of(&mut self, key: &[u8]) -> Option<Kind> {
let now = self.clock.now_ms();
let (kind, dead) = self
.map
.get(key)
.map(|rec| (value::kind(rec), value::is_expired(rec, now)))?;
if dead {
self.drop_key(key);
self.expired += 1;
return None;
}
Some(kind)
}
pub fn set_encoding(&mut self, key: &[u8]) -> Option<set::Encoding> {
self.reap(key);
let rec = self.map.get(key)?;
if value::kind(rec) != Kind::Set {
return None;
}
let at = value::slot(rec);
Some(self.sets.get(at)?.encoding())
}
pub fn hash_encoding(&mut self, key: &[u8]) -> Option<hash::Encoding> {
self.reap(key);
let rec = self.map.get(key)?;
if value::kind(rec) != Kind::Hash {
return None;
}
let at = value::slot(rec);
Some(self.hashes.get(at)?.encoding())
}
pub fn list_encoding(&mut self, key: &[u8]) -> Option<list::Encoding> {
self.reap(key);
let rec = self.map.get(key)?;
if value::kind(rec) != Kind::List {
return None;
}
let at = value::slot(rec);
Some(self.lists.get(at)?.encoding())
}
pub fn zset_encoding(&mut self, key: &[u8]) -> Option<zset::Encoding> {
self.reap(key);
let rec = self.map.get(key)?;
if value::kind(rec) != Kind::Zset {
return None;
}
let at = value::slot(rec);
Some(self.zsets.get(at)?.encoding())
}
pub fn encoding_name(&mut self, key: &[u8]) -> Option<&'static str> {
match self.kind_of(key)? {
Kind::String => self.encoding(key).map(value::Encoding::name),
Kind::Set => self.set_encoding(key).map(set::Encoding::name),
Kind::Hash => self.hash_encoding(key).map(hash::Encoding::name),
Kind::List => self.list_encoding(key).map(list::Encoding::name),
Kind::Zset => self.zset_encoding(key).map(zset::Encoding::name),
Kind::Array => Some("sliced-array"),
Kind::Stream => unreachable!("nothing can store a stream yet"),
}
}
pub fn set_expiry(&mut self, key: &[u8], at: Option<u64>) -> bool {
self.reap(key);
let Some(rec) = self.map.get(key) else {
return false;
};
if value::expire_at(rec) == at {
return true;
}
match value::kind(rec) {
Kind::String => {
let mut bytes = std::mem::take(&mut self.scratch);
bytes.clear();
value::read(rec).write_to(&mut bytes);
self.store(key, &bytes, at);
self.scratch = bytes;
}
kind @ (Kind::Set | Kind::Hash | Kind::List | Kind::Zset | Kind::Array) => {
let slot = value::slot(rec);
let len = value::slot_record_len(at.is_some());
self.map.set_with(key, len, |out| {
value::write_slot_record(out, kind, slot, at);
});
}
Kind::Stream => unreachable!("nothing can store a stream yet"),
}
true
}
pub fn deadline_of(&mut self, key: &[u8]) -> Ask {
let Some(addr) = self.live_rec(key) else {
return Ask::Missing;
};
match value::expire_at(self.map.value_at(addr)) {
Some(at) => Ask::At(at),
None => Ask::NoDeadline,
}
}
pub fn expire(&mut self, key: &[u8], at: u64, cond: Cond) -> Applied {
let prev = match self.deadline_of(key) {
Ask::Missing => return Applied::Missing,
Ask::NoDeadline => None,
Ask::At(at) => Some(at),
};
let done = ttl::decide(prev, at, cond, self.clock.now_ms());
match done {
Applied::Ok => {
self.set_expiry(key, Some(at));
}
Applied::Deleted => {
self.drop_key(key);
}
Applied::Missing | Applied::NotMet => {}
}
done
}
pub fn persist(&mut self, key: &[u8]) -> bool {
if !matches!(self.deadline_of(key), Ask::At(_)) {
return false;
}
self.set_expiry(key, None);
true
}
pub(crate) fn free_body(&mut self, key: &[u8]) {
if self.bodies == 0 {
return;
}
let Some(rec) = self.map.get(key) else {
return;
};
match value::kind(rec) {
Kind::String => {}
Kind::Set => {
let at = value::slot(rec);
self.sets.remove(at);
self.bodies -= 1;
}
Kind::Hash => {
let at = value::slot(rec);
self.hashes.remove(at);
self.bodies -= 1;
}
Kind::List => {
let at = value::slot(rec);
self.lists.remove(at);
self.bodies -= 1;
}
Kind::Zset => {
let at = value::slot(rec);
self.zsets.remove(at);
self.bodies -= 1;
}
Kind::Array => {
let at = value::slot(rec);
self.arrays.remove(at);
self.bodies -= 1;
}
Kind::Stream => unreachable!("nothing can store a stream yet"),
}
}
#[inline]
pub(crate) fn drop_key(&mut self, key: &[u8]) -> bool {
self.free_body(key);
self.map.del(key)
}
#[inline]
pub(crate) fn reap(&mut self, key: &[u8]) {
let now = self.clock.now_ms();
let dead = self.map.get(key).is_some_and(|r| value::is_expired(r, now));
if dead {
self.drop_key(key);
self.expired += 1;
}
}
pub(crate) fn live_rec(&mut self, key: &[u8]) -> Option<Addr> {
let now = self.clock.now_ms();
let addr = self.map.find(key)?;
if value::is_expired(self.map.value_at(addr), now) {
self.drop_key(key);
self.expired += 1;
return None;
}
Some(addr)
}
pub(crate) fn live_slot(&mut self, key: &[u8], want: Kind) -> Result<Option<u32>> {
if let Some((kind, slot)) = self.memo.get(self.map.writes(), key) {
if kind != want {
return Err(wrong_type());
}
return Ok(Some(slot));
}
let now = self.clock.now_ms();
let Some(rec) = self.map.get(key) else {
return Ok(None);
};
if value::is_expired(rec, now) {
self.drop_key(key);
self.expired += 1;
return Ok(None);
}
if value::kind(rec) != want {
return Err(wrong_type());
}
let slot = value::slot(rec);
let dated = value::expire_at(rec).is_some();
if !dated {
self.memo.put(self.map.writes(), key, want, slot);
}
Ok(Some(slot))
}
pub(crate) fn live_slot_either(
&mut self,
key: &[u8],
a: Kind,
b: Kind,
) -> Result<Option<(Kind, u32)>> {
if let Some((kind, slot)) = self.memo.get(self.map.writes(), key) {
if kind != a && kind != b {
return Err(wrong_type());
}
return Ok(Some((kind, slot)));
}
let now = self.clock.now_ms();
let Some(rec) = self.map.get(key) else {
return Ok(None);
};
if value::is_expired(rec, now) {
self.drop_key(key);
self.expired += 1;
return Ok(None);
}
let kind = value::kind(rec);
if kind != a && kind != b {
return Err(wrong_type());
}
let slot = value::slot(rec);
if value::expire_at(rec).is_none() {
self.memo.put(self.map.writes(), key, kind, slot);
}
Ok(Some((kind, slot)))
}
pub fn clear(&mut self) {
self.map.clear();
self.sets.clear();
self.hashes.clear();
self.lists.clear();
self.zsets.clear();
self.arrays.clear();
self.bodies = 0;
}
#[inline]
pub const fn expired_keys(&self) -> u64 {
self.expired
}
#[inline]
pub fn memory_bytes(&self) -> usize {
self.map.memory_bytes()
+ self.sets.memory_bytes()
+ self.sets.iter().map(Set::memory_bytes).sum::<usize>()
+ self.hashes.memory_bytes()
+ self.hashes.iter().map(Hash::memory_bytes).sum::<usize>()
+ self.lists.memory_bytes()
+ self.lists.iter().map(List::memory_bytes).sum::<usize>()
+ self.zsets.memory_bytes()
+ self.zsets.iter().map(Zset::memory_bytes).sum::<usize>()
+ self.arrays.memory_bytes()
+ self.arrays.iter().map(Array::memory_bytes).sum::<usize>()
}
#[inline]
pub fn compact_step(&mut self) -> Option<usize> {
self.map.compact_step()
}
#[inline]
pub fn prefetch(&self, hash: u64) {
self.map.prefetch(hash);
}
#[inline]
#[must_use]
pub fn hash_of(key: &[u8]) -> u64 {
RawMap::hash_of(key)
}
}
pub fn wrong_type() -> Error {
Error::new(
Code::WrongType,
"Operation against a key holding the wrong kind of value",
)
}
impl Default for Keyspace {
fn default() -> Keyspace {
Keyspace::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn db() -> Keyspace {
Keyspace::with_clock(Clock::fixed(1_000))
}
#[test]
fn type_answers_string_for_a_string_and_nothing_for_a_missing_key() {
let mut d = db();
d.set_plain(b"k", b"v").expect("room");
assert_eq!(d.kind_of(b"k"), Some(Kind::String));
assert_eq!(d.kind_of(b"nope"), None);
}
#[test]
fn type_does_not_report_a_key_whose_deadline_has_gone() {
let mut d = db();
d.psetex(b"k", 100, b"v").expect("room");
assert_eq!(d.kind_of(b"k"), Some(Kind::String));
d.clock_mut().advance(100);
assert_eq!(
d.kind_of(b"k"),
None,
"the deadline was 1100 and it is 1100"
);
assert_eq!(d.len(), 0, "and asking reaped it rather than leaving it");
assert_eq!(d.expired_keys(), 1);
}
}