use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
use std::sync::{Arc, Condvar, Mutex, RwLock};
use std::time::Instant;
use crate::codec::{BincodeCodec, Codec};
use crate::entry::{
DequeEntry, Entry, HashEntry, ListEntry, ListInner, Meta, SetEntry, StringEntry, ValueType,
ZSetEntry,
};
use crate::error::Error;
pub struct Store<C: Codec = BincodeCodec> {
codec: C,
map: RwLock<HashMap<String, Entry>>,
}
impl Default for Store<BincodeCodec> {
fn default() -> Self {
Self::new()
}
}
impl Store<BincodeCodec> {
pub fn new() -> Self {
Self::with_codec(BincodeCodec)
}
}
impl<C: Codec> Store<C> {
pub fn with_codec(codec: C) -> Self {
Self {
codec,
map: RwLock::new(HashMap::new()),
}
}
pub fn codec(&self) -> &C {
&self.codec
}
pub fn kv(&self) -> crate::kv::KvOps<'_, C> {
crate::kv::KvOps::new(self)
}
pub fn keys(&self) -> crate::keys::KeysOps<'_, C> {
crate::keys::KeysOps::new(self)
}
pub fn hash<'a>(&'a self, key: &'a str) -> crate::hash::HashRef<'a, C> {
crate::hash::HashRef::new(self, key)
}
pub fn set<'a>(&'a self, key: &'a str) -> crate::set::SetRef<'a, C> {
crate::set::SetRef::new(self, key)
}
pub fn list<'a>(&'a self, key: &'a str) -> crate::list::ListRef<'a, C> {
crate::list::ListRef::new(self, key)
}
pub fn deque<'a>(&'a self, key: &'a str) -> crate::deque::DequeRef<'a, C> {
crate::deque::DequeRef::new(self, key)
}
pub fn zset<'a>(&'a self, key: &'a str) -> crate::zset::ZSetRef<'a, C> {
crate::zset::ZSetRef::new(self, key)
}
pub(crate) fn with_map_read<R>(&self, f: impl FnOnce(&HashMap<String, Entry>) -> R) -> R {
let guard = self.map.read().expect("store lock poisoned");
f(&guard)
}
pub(crate) fn with_map_write<R>(&self, f: impl FnOnce(&mut HashMap<String, Entry>) -> R) -> R {
let mut guard = self.map.write().expect("store lock poisoned");
f(&mut guard)
}
pub(crate) fn contains_key(&self, key: &str) -> bool {
self.with_map_read(|m| m.contains_key(key))
}
pub(crate) fn snapshot_keys(&self) -> Vec<String> {
self.with_map_read(|m| m.keys().cloned().collect())
}
pub(crate) fn get_entry(&self, key: &str) -> Option<Entry> {
self.with_map_read(|m| m.get(key).cloned())
}
pub(crate) fn get_entry_meta(&self, key: &str) -> Option<Meta> {
self.with_map_read(|m| m.get(key).map(|e| *e.meta()))
}
pub(crate) fn put_string_entry(&self, key: &str, bytes: crate::codec::Bytes) {
self.with_map_write(|m| {
let entry = Entry::String(StringEntry {
meta: Meta::new(ValueType::String),
bytes,
});
m.insert(key.to_string(), entry);
});
}
pub(crate) fn remove_entry(&self, key: &str) -> bool {
self.with_map_write(|m| m.remove(key).is_some())
}
pub(crate) fn set_expire_at(&self, key: &str, when: Option<Instant>) -> bool {
self.with_map_write(|m| {
let Some(entry) = m.get_mut(key) else {
return false;
};
entry.meta_mut().expire_at = when;
true
})
}
pub(crate) fn purge_if_expired(&self, key: &str) {
let expired = self.with_map_read(|m| {
let Some(entry) = m.get(key) else {
return false;
};
match entry.meta().expire_at {
None => false,
Some(t) => t <= Instant::now(),
}
});
if !expired {
return;
}
self.with_map_write(|m| {
let now = Instant::now();
let should_remove = match m.get(key) {
None => false,
Some(entry) => match entry.meta().expire_at {
None => false,
Some(t) => t <= now,
},
};
if should_remove {
m.remove(key);
}
});
}
pub(crate) fn rename_internal(&self, from: &str, to: &str, nx: bool) -> Result<bool, Error> {
if from == to {
return Ok(true);
}
self.with_map_write(|m| {
let Some(entry) = m.remove(from) else {
return Err(Error::NotFound);
};
if nx && m.contains_key(to) {
m.insert(from.to_string(), entry);
return Ok(false);
}
m.insert(to.to_string(), entry);
Ok(true)
})
}
pub(crate) fn with_hash_mut<R>(
&self,
key: &str,
f: impl FnOnce(&mut HashMap<crate::codec::Bytes, crate::codec::Bytes>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_write(|m| {
if !m.contains_key(key) {
m.insert(
key.to_string(),
Entry::Hash(HashEntry {
meta: Meta::new(ValueType::Hash),
map: HashMap::new(),
}),
);
}
match m.get_mut(key).expect("just inserted or existed") {
Entry::Hash(he) => f(&mut he.map),
other => Err(Error::WrongType {
expected: ValueType::Hash.as_str(),
got: other.value_type().as_str(),
}),
}
})
}
pub(crate) fn with_hash_read<R>(
&self,
key: &str,
f: impl FnOnce(Option<&HashMap<crate::codec::Bytes, crate::codec::Bytes>>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_read(|m| match m.get(key) {
None => f(None),
Some(entry) => match entry {
Entry::Hash(he) => f(Some(&he.map)),
other => Err(Error::WrongType {
expected: ValueType::Hash.as_str(),
got: other.value_type().as_str(),
}),
},
})
}
pub(crate) fn with_set_mut<R>(
&self,
key: &str,
f: impl FnOnce(&mut HashSet<crate::codec::Bytes>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_write(|m| {
if !m.contains_key(key) {
m.insert(
key.to_string(),
Entry::Set(SetEntry {
meta: Meta::new(ValueType::Set),
set: HashSet::new(),
}),
);
}
match m.get_mut(key).expect("just inserted or existed") {
Entry::Set(se) => f(&mut se.set),
other => Err(Error::WrongType {
expected: ValueType::Set.as_str(),
got: other.value_type().as_str(),
}),
}
})
}
pub(crate) fn with_set_read<R>(
&self,
key: &str,
f: impl FnOnce(Option<&HashSet<crate::codec::Bytes>>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_read(|m| match m.get(key) {
None => f(None),
Some(entry) => match entry {
Entry::Set(se) => f(Some(&se.set)),
other => Err(Error::WrongType {
expected: ValueType::Set.as_str(),
got: other.value_type().as_str(),
}),
},
})
}
pub(crate) fn with_list_mut<R>(
&self,
key: &str,
f: impl FnOnce(&mut VecDeque<crate::codec::Bytes>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
let inner = self.with_map_write(|m| {
match m.get(key) {
None => {
let entry = Entry::List(ListEntry {
meta: Meta::new(ValueType::List),
inner: Arc::new((
Mutex::new(ListInner {
deque: VecDeque::new(),
}),
Condvar::new(),
)),
});
m.insert(key.to_string(), entry);
}
Some(e) if e.value_type() != ValueType::List => {
return Err(Error::WrongType {
expected: ValueType::List.as_str(),
got: e.value_type().as_str(),
});
}
Some(_) => {}
}
match m.get(key).expect("ensured above") {
Entry::List(le) => Ok(le.inner.clone()),
other => Err(Error::WrongType {
expected: ValueType::List.as_str(),
got: other.value_type().as_str(),
}),
}
})?;
let (lock, _cv) = &*inner;
let mut guard = lock.lock().expect("list mutex poisoned");
f(&mut guard.deque)
}
pub(crate) fn with_list_read<R>(
&self,
key: &str,
f: impl FnOnce(Option<&VecDeque<crate::codec::Bytes>>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
let inner = self.with_map_read(|m| match m.get(key) {
None => Ok::<_, Error>(None),
Some(entry) => match entry {
Entry::List(le) => Ok(Some(le.inner.clone())),
other => Err(Error::WrongType {
expected: ValueType::List.as_str(),
got: other.value_type().as_str(),
}),
},
})?;
let Some(inner) = inner else {
return f(None);
};
let (lock, _cv) = &*inner;
let guard = lock.lock().expect("list mutex poisoned");
f(Some(&guard.deque))
}
pub(crate) fn list_notify(&self, key: &str) -> Result<(), Error> {
self.purge_if_expired(key);
let inner = self.with_map_read(|m| match m.get(key) {
None => Ok::<_, Error>(None),
Some(entry) => match entry {
Entry::List(le) => Ok(Some(le.inner.clone())),
other => Err(Error::WrongType {
expected: ValueType::List.as_str(),
got: other.value_type().as_str(),
}),
},
})?;
if let Some(inner) = inner {
let (_lock, cv) = &*inner;
cv.notify_all();
}
Ok(())
}
pub(crate) fn with_deque_mut<R>(
&self,
key: &str,
f: impl FnOnce(&mut VecDeque<crate::codec::Bytes>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_write(|m| {
let entry = m.entry(key.to_string()).or_insert_with(|| {
Entry::Deque(DequeEntry {
meta: Meta::new(ValueType::Deque),
deque: VecDeque::new(),
})
});
match entry {
Entry::Deque(de) => f(&mut de.deque),
other => Err(Error::WrongType {
expected: ValueType::Deque.as_str(),
got: other.value_type().as_str(),
}),
}
})
}
pub(crate) fn with_deque_read<R>(
&self,
key: &str,
f: impl FnOnce(Option<&VecDeque<crate::codec::Bytes>>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_read(|m| match m.get(key) {
None => f(None),
Some(entry) => match entry {
Entry::Deque(de) => f(Some(&de.deque)),
other => Err(Error::WrongType {
expected: ValueType::Deque.as_str(),
got: other.value_type().as_str(),
}),
},
})
}
pub(crate) fn with_zset_mut<R>(
&self,
key: &str,
f: impl FnOnce(&mut ZSetEntry) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_write(|m| {
let entry = m.entry(key.to_string()).or_insert_with(|| {
Entry::ZSet(ZSetEntry {
meta: Meta::new(ValueType::ZSet),
member_to_score: HashMap::new(),
score_to_members: BTreeMap::new(),
})
});
match entry {
Entry::ZSet(ze) => f(ze),
other => Err(Error::WrongType {
expected: ValueType::ZSet.as_str(),
got: other.value_type().as_str(),
}),
}
})
}
pub(crate) fn with_zset_read<R>(
&self,
key: &str,
f: impl FnOnce(Option<&ZSetEntry>) -> Result<R, Error>,
) -> Result<R, Error> {
self.purge_if_expired(key);
self.with_map_read(|m| match m.get(key) {
None => f(None),
Some(entry) => match entry {
Entry::ZSet(ze) => f(Some(ze)),
other => Err(Error::WrongType {
expected: ValueType::ZSet.as_str(),
got: other.value_type().as_str(),
}),
},
})
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use crate::Store;
#[test]
fn rename_preserves_structure() {
let store = Store::new();
let h = store.hash("old");
h.hset("a", &1_i64).unwrap();
store.keys().rename("old", "new").unwrap();
let h2 = store.hash("new");
let v: Option<i64> = h2.hget("a").unwrap();
assert_eq!(v, Some(1));
}
#[test]
fn ttl_applies_to_hash() {
let store = Store::new();
let h = store.hash("h");
h.hset("a", &1_i64).unwrap();
store.keys().pexpire("h", 1);
std::thread::sleep(Duration::from_millis(3));
let v: Option<i64> = h.hget("a").unwrap();
assert_eq!(v, None);
assert!(!store.keys().exists("h"));
}
}