pub mod conf;
pub mod meta;
pub use conf::{
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING, ERR_HASH_VALUE_NOT_FLOAT,
ERR_HASH_VALUE_NOT_INTEGER, ERR_INCREMENT_NAN_OR_INFINITY, ERR_INCREMENT_OVERFLOW, FieldValue,
HASH_EXPIRE_COND_FAILED, HASH_EXPIRE_DELETED, HASH_EXPIRE_SET_OK, HASH_FIELD_NOT_FOUND,
HASH_FIELD_PERSISTENT, HExpire, HGetEx, HSetEx, HashFetchType, HashFieldExpireCondition,
HashFieldSetCondition, HashGetExOptions, HashLengthMode, HashSetExOptions, RangeLexSpec,
TTLAction,
};
pub use meta::{
FIELD_EXPIRE_PREFIX_LEN, HashFieldState, HashFieldStateKind, HashMeta, HashSubkeyEncodingMode,
decode_field_state, decode_hash_value, encode_hash_value, hexpire_condition_passes,
is_field_expired, is_immediate_expire,
};
pub type HashFieldPair = (Vec<u8>, Vec<u8>);
pub type HashRandField = (Vec<u8>, Option<Vec<u8>>);
pub type HashScanResult = (usize, Vec<HashFieldPair>);
use rapidhash::{RapidHashMap as HashMap, RapidHashSet as HashSet};
use std::str;
use crate::db::WeDb;
use crate::error::{Error, Result};
use crate::key_composer::{KeyComposer, matches_glob_bytes};
#[inline]
fn current_now_ms() -> u64 {
jiff::Timestamp::now().as_millisecond().max(0) as u64
}
#[inline]
const fn ceil_div_1000(val: u64) -> u64 {
val / 1000 + if !val.is_multiple_of(1000) { 1 } else { 0 }
}
#[derive(Clone)]
struct CachedFieldState {
kind: HashFieldStateKind,
expire: u64,
raw: Option<fjall::Slice>,
}
impl WeDb {
#[inline]
fn get_hash_meta(&self, meta_k: &str, now_ms: u64) -> Result<Option<HashMeta>> {
match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => {
if let Some(meta) = HashMeta::decode(&m_bytes)
&& !meta.is_expired(now_ms)
{
return Ok(Some(meta));
}
Ok(None)
}
None => Ok(None),
}
}
#[inline]
fn prepare_hash_meta_for_write(
&self,
kc: &KeyComposer,
k_str: &str,
meta_k: &str,
now_ms: u64,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<(HashMeta, bool)> {
match self.meta_ks.get(meta_k.as_bytes())? {
Some(m_bytes) => match HashMeta::decode(&m_bytes) {
Some(meta) if !meta.is_expired(now_ms) => Ok((meta, true)),
Some(_) => {
let prefix = kc.hash_prefix(k_str);
for g in self.data_ks.prefix(&prefix) {
let (k, _) = g.into_inner()?;
if k.starts_with(&prefix) {
batch.remove(&self.data_ks, &*k);
} else {
break;
}
}
Ok((HashMeta::new_with_version(0, 0), false))
}
None => Ok((HashMeta::new_with_version(0, 0), false)),
},
None => Ok((HashMeta::new_with_version(0, 0), false)),
}
}
pub fn hset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_vals: &[(F, V)],
) -> Result<usize> {
if field_vals.is_empty() {
return Ok(0);
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let (mut meta, metadata_existed) =
self.prepare_hash_meta_for_write(&kc, k_str, &meta_k, now_ms, &mut batch)?;
let mut seen = HashSet::with_capacity_and_hasher(field_vals.len(), Default::default());
let mut unique_items = Vec::with_capacity(field_vals.len());
for (f, v) in field_vals.iter().rev() {
let f_bytes = f.as_ref();
if seen.insert(f_bytes) {
unique_items.push((f_bytes, v.as_ref()));
}
}
unique_items.reverse();
let mut added = 0;
let mut meta_changed = false;
for (f_bytes, v_bytes) in unique_items {
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let kind = if metadata_existed
&& let Some(raw) = self.data_ks.get(&item_k)?
&& let Some(state) = decode_field_state(&meta, &raw, now_ms)
{
if state.kind == HashFieldStateKind::Persistent && state.value == v_bytes {
continue;
}
state.kind
} else {
HashFieldStateKind::Missing
};
match kind {
HashFieldStateKind::Missing => {
added += 1;
meta.apply_missing_to_persistent();
meta_changed = true;
}
HashFieldStateKind::Persistent => {
}
HashFieldStateKind::LiveTTL => {
meta.apply_ttl_to_persistent();
meta_changed = true;
}
HashFieldStateKind::ExpiredTTLPhysical => {
added += 1;
meta.apply_ttl_to_persistent();
meta_changed = true;
}
}
let enc_val = meta.encode_subkey_value(v_bytes, 0);
batch.insert(&self.data_ks, &item_k, enc_val);
}
if meta_changed || !metadata_existed {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(added)
}
#[inline]
pub fn hmset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_vals: &[(F, V)],
) -> Result<()> {
self.hset(key, field_vals)?;
Ok(())
}
pub fn hsetnx<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field: F,
value: V,
) -> Result<bool> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let (mut meta, metadata_existed) =
self.prepare_hash_meta_for_write(&kc, k_str, &meta_k, now_ms, &mut batch)?;
let f_bytes = field.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let kind = if metadata_existed
&& let Some(raw) = self.data_ks.get(&item_k)?
&& let Some(state) = decode_field_state(&meta, &raw, now_ms)
{
state.kind
} else {
HashFieldStateKind::Missing
};
let should_insert = match kind {
HashFieldStateKind::Missing => {
meta.apply_missing_to_persistent();
true
}
HashFieldStateKind::Persistent | HashFieldStateKind::LiveTTL => false,
HashFieldStateKind::ExpiredTTLPhysical => {
meta.apply_ttl_to_persistent();
true
}
};
if !should_insert {
return Ok(false);
}
let enc_val = meta.encode_subkey_value(value.as_ref(), 0);
batch.insert(&self.data_ks, &item_k, enc_val);
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
batch.commit()?;
Ok(true)
}
pub fn hget<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
) -> Result<Option<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(None),
};
let item_k = kc.hash_key_bytes(k_str, field.as_ref());
match self.data_ks.get(&item_k)? {
Some(raw) => match meta.decode_subkey_value(&raw) {
Some((exp, payload)) => {
if is_field_expired(exp, now_ms) {
Ok(None)
} else {
Ok(Some(payload.to_vec()))
}
}
None => Ok(None),
},
None => Ok(None),
}
}
pub fn hdel<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<usize> {
if fields.is_empty() {
return Ok(0);
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
let mut seen = HashSet::with_capacity_and_hasher(fields.len(), Default::default());
let mut deleted = 0;
let mut physical_removed = 0;
let mut batch = self.db.batch();
for f in fields {
let f_bytes = f.as_ref();
if !seen.insert(f_bytes) {
continue;
}
let item_k = kc.hash_key_bytes(k_str, f_bytes);
if let Some(raw) = self.data_ks.get(&item_k)?
&& let Some(state) = decode_field_state(&meta, &raw, now_ms)
{
match state.kind {
HashFieldStateKind::Persistent => {
deleted += 1;
physical_removed += 1;
meta.apply_persistent_to_deleted();
batch.remove(&self.data_ks, &item_k);
}
HashFieldStateKind::LiveTTL => {
deleted += 1;
physical_removed += 1;
meta.apply_ttl_to_deleted();
batch.remove(&self.data_ks, &item_k);
}
HashFieldStateKind::ExpiredTTLPhysical => {
physical_removed += 1;
meta.apply_ttl_to_deleted();
batch.remove(&self.data_ks, &item_k);
}
HashFieldStateKind::Missing => {}
}
}
}
if physical_removed == 0 {
return Ok(0);
}
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(deleted)
}
#[inline]
pub fn hexists<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<bool> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(false),
};
let item_k = kc.hash_key_bytes(k_str, field.as_ref());
if let Some(raw) = self.data_ks.get(&item_k)?
&& let Some((exp, _)) = meta.decode_subkey_value(&raw)
{
return Ok(!is_field_expired(exp, now_ms));
}
Ok(false)
}
#[inline]
pub fn hlen<K: AsRef<[u8]>>(&self, key: K) -> Result<usize> {
self.hlen_with_mode(key, HashLengthMode::Accurate)
}
pub fn hlen_with_mode<K: AsRef<[u8]>>(&self, key: K, mode: HashLengthMode) -> Result<usize> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
if meta.is_legacy_subkey_encoding() || mode == HashLengthMode::Approximate {
return Ok(meta.base.size as usize);
}
if meta.persist > meta.base.size {
return self.scan_and_repair_hash(&kc, k_str, &mut meta, now_ms);
}
let ttl_candidates = meta.base.size.saturating_sub(meta.persist);
if ttl_candidates == 0 {
return Ok(meta.base.size as usize);
}
if meta.lower != 0 && now_ms < meta.lower {
return Ok(meta.base.size as usize);
}
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
let mut batch = self.db.batch();
batch.remove(&self.meta_ks, meta_k.as_bytes());
batch.commit()?;
return Ok(0);
}
self.scan_and_repair_hash(&kc, k_str, &mut meta, now_ms)
}
fn scan_and_repair_hash(
&self,
kc: &KeyComposer,
k_str: &str,
meta: &mut HashMeta,
now_ms: u64,
) -> Result<usize> {
let meta_k = kc.hash_meta(k_str);
let prefix = kc.hash_prefix(k_str);
let mut repaired = *meta;
repaired.base.size = 0;
repaired.persist = 0;
repaired.lower = 0;
repaired.upper = 0;
let mut batch = self.db.batch();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some(state) = decode_field_state(meta, &v, now_ms) {
match state.kind {
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(&self.data_ks, &*k);
}
HashFieldStateKind::Persistent => {
repaired.base.size += 1;
repaired.persist += 1;
}
HashFieldStateKind::LiveTTL => {
repaired.base.size += 1;
if repaired.lower == 0 || state.expire < repaired.lower {
repaired.lower = state.expire;
}
repaired.upper = repaired.upper.max(state.expire);
}
HashFieldStateKind::Missing => {}
}
}
}
if repaired.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
repaired.clear_bounds_if_no_ttl_candidates();
batch.insert(&self.meta_ks, meta_k.as_bytes(), repaired.encode());
}
batch.commit()?;
*meta = repaired;
Ok(repaired.base.size as usize)
}
pub fn hmget<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![None; fields.len()]),
};
let mut results = Vec::with_capacity(fields.len());
for f in fields {
let item_k = kc.hash_key_bytes(k_str, f.as_ref());
let val = match self.data_ks.get(&item_k)? {
Some(raw) => match meta.decode_subkey_value(&raw) {
Some((exp, payload)) => {
if is_field_expired(exp, now_ms) {
None
} else {
Some(payload.to_vec())
}
}
None => None,
},
None => None,
};
results.push(val);
}
Ok(results)
}
pub fn hget_all_with_type<K: AsRef<[u8]>>(
&self,
key: K,
fetch_type: HashFetchType,
) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
let prefix = kc.hash_prefix(k_str);
let mut items = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
let field_bytes = &k[prefix.len()..];
if let Some((exp, payload)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
match fetch_type {
HashFetchType::All => {
items.push((field_bytes.to_vec(), payload.to_vec()));
}
HashFetchType::OnlyKey => {
items.push((field_bytes.to_vec(), Vec::new()));
}
HashFetchType::OnlyValue => {
items.push((Vec::new(), payload.to_vec()));
}
}
}
}
Ok(items)
}
#[inline]
pub fn hgetall<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
self.hget_all_with_type(key, HashFetchType::All)
}
pub fn hkeys<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
let prefix = kc.hash_prefix(k_str);
let mut keys = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some((exp, _)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
keys.push(k[prefix.len()..].to_vec());
}
}
Ok(keys)
}
pub fn hvals<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<Vec<u8>>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
let prefix = kc.hash_prefix(k_str);
let mut vals = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
if let Some((exp, payload)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
vals.push(payload.to_vec());
}
}
Ok(vals)
}
pub fn hincrby<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
increment: i64,
) -> Result<i64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let (mut meta, metadata_existed) =
self.prepare_hash_meta_for_write(&kc, k_str, &meta_k, now_ms, &mut batch)?;
let f_bytes = field.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let mut cur_num: i64 = 0;
let mut keep_expire: u64 = 0;
let mut is_missing = false;
let mut is_expired_ttl = false;
if metadata_existed && let Some(raw) = self.data_ks.get(&item_k)? {
if let Some((exp, payload)) = meta.decode_subkey_value(&raw) {
if exp > 0 && is_field_expired(exp, now_ms) {
is_expired_ttl = true;
} else {
let s = str::from_utf8(payload)
.map_err(|_| Error::invalid_data(ERR_HASH_VALUE_NOT_INTEGER))?;
cur_num = s
.parse::<i64>()
.map_err(|_| Error::invalid_data(ERR_HASH_VALUE_NOT_INTEGER))?;
keep_expire = exp;
}
} else {
is_missing = true;
}
} else {
is_missing = true;
}
let new_num = cur_num
.checked_add(increment)
.ok_or_else(|| Error::invalid_data(ERR_INCREMENT_OVERFLOW))?;
let mut meta_changed = false;
if is_missing {
meta.apply_missing_to_persistent();
meta_changed = true;
} else if is_expired_ttl {
meta.apply_ttl_to_persistent();
meta_changed = true;
}
let mut buf = itoa::Buffer::new();
let num_str = buf.format(new_num);
let enc_val = meta.encode_subkey_value(num_str.as_bytes(), keep_expire);
batch.insert(&self.data_ks, &item_k, enc_val);
if meta_changed || !metadata_existed {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(new_num)
}
pub fn hincrbyfloat<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
increment: f64,
) -> Result<f64> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let (mut meta, metadata_existed) =
self.prepare_hash_meta_for_write(&kc, k_str, &meta_k, now_ms, &mut batch)?;
let f_bytes = field.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let mut cur_num: f64 = 0.0;
let mut keep_expire: u64 = 0;
let mut is_missing = false;
let mut is_expired_ttl = false;
if metadata_existed && let Some(raw) = self.data_ks.get(&item_k)? {
if let Some((exp, payload)) = meta.decode_subkey_value(&raw) {
if exp > 0 && is_field_expired(exp, now_ms) {
is_expired_ttl = true;
} else {
let s = str::from_utf8(payload)
.map_err(|_| Error::invalid_data(ERR_HASH_VALUE_NOT_FLOAT))?;
cur_num = s
.parse::<f64>()
.map_err(|_| Error::invalid_data(ERR_HASH_VALUE_NOT_FLOAT))?;
keep_expire = exp;
}
} else {
is_missing = true;
}
} else {
is_missing = true;
}
let new_num = cur_num + increment;
if new_num.is_nan() || new_num.is_infinite() {
return Err(Error::invalid_data(ERR_INCREMENT_NAN_OR_INFINITY));
}
let mut meta_changed = false;
if is_missing {
meta.apply_missing_to_persistent();
meta_changed = true;
} else if is_expired_ttl {
meta.apply_ttl_to_persistent();
meta_changed = true;
}
let mut buf = ryu::Buffer::new();
let num_str = buf.format(new_num);
let enc_val = meta.encode_subkey_value(num_str.as_bytes(), keep_expire);
batch.insert(&self.data_ks, &item_k, enc_val);
if meta_changed || !metadata_existed {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(new_num)
}
#[inline]
pub fn hstrlen<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<usize> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(0),
};
let item_k = kc.hash_key_bytes(k_str, field.as_ref());
if let Some(raw) = self.data_ks.get(&item_k)?
&& let Some((exp, payload)) = meta.decode_subkey_value(&raw)
&& !is_field_expired(exp, now_ms)
{
return Ok(payload.len());
}
Ok(0)
}
pub fn hexpire<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
seconds: i64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
let target_expire_ms = if seconds <= 0 {
0
} else {
now_ms.saturating_add((seconds as u64).saturating_mul(1000))
};
self.expire_fields(key, fields, target_expire_ms, condition, now_ms)
}
pub fn hpexpire<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
milliseconds: i64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
let target_expire_ms = if milliseconds <= 0 {
0
} else {
now_ms.saturating_add(milliseconds as u64)
};
self.expire_fields(key, fields, target_expire_ms, condition, now_ms)
}
pub fn hexpireat<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
unix_time_sec: u64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
let target_expire_ms = unix_time_sec.saturating_mul(1000);
self.expire_fields(key, fields, target_expire_ms, condition, now_ms)
}
pub fn hpexpireat<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
unix_time_ms: u64,
condition: HExpire,
) -> Result<Vec<i64>> {
let now_ms = current_now_ms();
self.expire_fields(key, fields, unix_time_ms, condition, now_ms)
}
pub fn expire_fields<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
expire_at_ms: u64,
condition: HExpire,
now_ms: u64,
) -> Result<Vec<i64>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let mut meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let is_immediate = is_immediate_expire(expire_at_ms, now_ms);
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.db.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
for f in fields {
let f_bytes = f.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let raw_opt = self.data_ks.get(&item_k)?;
let state_entry = match raw_opt {
Some(raw) => {
if let Some(state) = decode_field_state(&meta, &raw, now_ms) {
CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
}
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
}
}
None => CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
};
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(&self.data_ks, &item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::Persistent | HashFieldStateKind::LiveTTL => {
if !hexpire_condition_passes(condition, entry.kind, entry.expire, expire_at_ms)
{
results.push(HASH_EXPIRE_COND_FAILED);
continue;
}
if is_immediate {
batch.remove(&self.data_ks, &item_k);
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_deleted();
} else {
meta.apply_ttl_to_deleted();
}
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(HASH_EXPIRE_DELETED);
} else {
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_ttl(expire_at_ms);
} else {
meta.apply_ttl_to_ttl(expire_at_ms);
}
meta_changed = true;
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p)
.unwrap_or(b"");
let enc = meta.encode_subkey_value(payload, expire_at_ms);
batch.insert(&self.data_ks, &item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::LiveTTL,
expire: expire_at_ms,
raw: entry.raw,
},
);
results.push(HASH_EXPIRE_SET_OK);
}
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
}
Ok(results)
}
pub fn httl<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
let mut result_cache: HashMap<&[u8], i64> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
let mut results = Vec::with_capacity(fields.len());
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let res = match self.data_ks.get(&item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => {
let remain_ms = s.expire.saturating_sub(now_ms);
ceil_div_1000(remain_ms) as i64
}
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
pub fn hpttl<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, fields: &[F]) -> Result<Vec<i64>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
let mut result_cache: HashMap<&[u8], i64> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
let mut results = Vec::with_capacity(fields.len());
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let res = match self.data_ks.get(&item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => {
let remain_ms = s.expire.saturating_sub(now_ms);
remain_ms as i64
}
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
pub fn hexpiretime<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<i64>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
let mut result_cache: HashMap<&[u8], i64> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
let mut results = Vec::with_capacity(fields.len());
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let res = match self.data_ks.get(&item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => ceil_div_1000(s.expire) as i64,
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
pub fn hpexpiretime<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<i64>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
let mut result_cache: HashMap<&[u8], i64> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
let mut results = Vec::with_capacity(fields.len());
for f in fields {
let f_bytes = f.as_ref();
if let Some(&cached) = result_cache.get(f_bytes) {
results.push(cached);
continue;
}
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let res = match self.data_ks.get(&item_k)? {
None => HASH_FIELD_NOT_FOUND,
Some(raw) => match decode_field_state(&meta, &raw, now_ms) {
None => HASH_FIELD_NOT_FOUND,
Some(s) => match s.kind {
HashFieldStateKind::Missing | HashFieldStateKind::ExpiredTTLPhysical => {
HASH_FIELD_NOT_FOUND
}
HashFieldStateKind::Persistent => HASH_FIELD_PERSISTENT,
HashFieldStateKind::LiveTTL => s.expire as i64,
},
},
};
result_cache.insert(f_bytes, res);
results.push(res);
}
Ok(results)
}
pub fn hpersist<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<i64>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![HASH_FIELD_NOT_FOUND; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.db.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
for f in fields {
let f_bytes = f.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let raw_opt = self.data_ks.get(&item_k)?;
let state_entry = match raw_opt {
Some(raw) => {
if let Some(state) = decode_field_state(&meta, &raw, now_ms) {
CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
}
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
}
}
None => CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
};
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::Persistent => {
results.push(HASH_FIELD_PERSISTENT);
}
HashFieldStateKind::ExpiredTTLPhysical => {
batch.remove(&self.data_ks, &item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
results.push(HASH_FIELD_NOT_FOUND);
}
HashFieldStateKind::LiveTTL => {
meta.apply_ttl_to_persistent();
meta_changed = true;
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p)
.unwrap_or(b"");
let enc = meta.encode_subkey_value(payload, 0);
batch.insert(&self.data_ks, &item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Persistent,
expire: 0,
raw: entry.raw,
},
);
results.push(HASH_EXPIRE_SET_OK);
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
}
Ok(results)
}
pub fn hgetdel<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![None; fields.len()]),
};
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.db.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
for f in fields {
let f_bytes = f.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let raw_opt = self.data_ks.get(&item_k)?;
let state_entry = match raw_opt {
Some(raw) => {
if let Some(state) = decode_field_state(&meta, &raw, now_ms) {
CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
}
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
}
}
None => CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
};
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(None);
}
HashFieldStateKind::Persistent => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p.to_vec())
.unwrap_or_default();
results.push(Some(payload));
batch.remove(&self.data_ks, &item_k);
meta.apply_persistent_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
HashFieldStateKind::LiveTTL => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p.to_vec())
.unwrap_or_default();
results.push(Some(payload));
batch.remove(&self.data_ks, &item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
HashFieldStateKind::ExpiredTTLPhysical => {
results.push(None);
batch.remove(&self.data_ks, &item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
}
Ok(results)
}
pub fn set_fields_with_expire<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_values: &[(F, V)],
options: HashSetExOptions,
) -> Result<bool> {
if field_values.is_empty() {
return Ok(false);
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut batch = self.db.batch();
let (mut meta, metadata_existed) =
self.prepare_hash_meta_for_write(&kc, k_str, &meta_k, now_ms, &mut batch)?;
if !metadata_existed && options.condition == HashFieldSetCondition::Fxx {
return Ok(false);
}
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> =
HashMap::with_capacity_and_hasher(field_values.len(), Default::default());
if options.condition != HashFieldSetCondition::None {
for (f, _) in field_values {
let f_bytes = f.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let raw_opt = if metadata_existed {
self.data_ks.get(&item_k)?
} else {
None
};
let state_entry = match raw_opt {
Some(raw) => {
if let Some(state) = decode_field_state(&meta, &raw, now_ms) {
CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
}
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
}
}
None => CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
};
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
if entry.kind == HashFieldStateKind::ExpiredTTLPhysical {
batch.remove(&self.data_ks, &item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
let exists = matches!(
entry.kind,
HashFieldStateKind::Persistent | HashFieldStateKind::LiveTTL
);
if (options.condition == HashFieldSetCondition::Fnx && exists)
|| (options.condition == HashFieldSetCondition::Fxx && !exists)
{
if meta_changed {
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
}
return Ok(false);
}
}
}
let mut seen = HashSet::with_capacity_and_hasher(field_values.len(), Default::default());
let mut unique_field_values = Vec::with_capacity(field_values.len());
for (f, v) in field_values.iter().rev() {
let f_bytes = f.as_ref();
if seen.insert(f_bytes) {
unique_field_values.push((f_bytes, v.as_ref()));
}
}
unique_field_values.reverse();
let is_immediate = options.ttl_action == TTLAction::Set
&& is_immediate_expire(options.expire_at_ms, now_ms);
for (f_bytes, v_bytes) in unique_field_values {
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let raw_opt = if metadata_existed {
self.data_ks.get(&item_k)?
} else {
None
};
let state_entry = match raw_opt {
Some(raw) => {
if let Some(state) = decode_field_state(&meta, &raw, now_ms) {
CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
}
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
}
}
None => CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
};
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
if is_immediate {
match entry.kind {
HashFieldStateKind::Missing => continue,
HashFieldStateKind::Persistent => {
meta.apply_persistent_to_deleted();
meta_changed = true;
}
HashFieldStateKind::LiveTTL | HashFieldStateKind::ExpiredTTLPhysical => {
meta.apply_ttl_to_deleted();
meta_changed = true;
}
}
batch.remove(&self.data_ks, &item_k);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
continue;
}
let target_expire = match options.ttl_action {
TTLAction::Discard | TTLAction::Persist => 0,
TTLAction::Keep => {
if entry.kind == HashFieldStateKind::LiveTTL
|| entry.kind == HashFieldStateKind::ExpiredTTLPhysical
{
entry.expire
} else {
0
}
}
TTLAction::Set => options.expire_at_ms,
};
match entry.kind {
HashFieldStateKind::Missing => {
if target_expire == 0 {
meta.apply_missing_to_persistent();
} else {
meta.apply_missing_to_ttl(target_expire);
}
meta_changed = true;
}
HashFieldStateKind::Persistent => {
if target_expire != 0 {
meta.apply_persistent_to_ttl(target_expire);
meta_changed = true;
}
}
HashFieldStateKind::LiveTTL | HashFieldStateKind::ExpiredTTLPhysical => {
if target_expire == 0 {
meta.apply_ttl_to_persistent();
} else {
meta.apply_ttl_to_ttl(target_expire);
}
meta_changed = true;
}
}
let enc_val = meta.encode_subkey_value(v_bytes, target_expire);
batch.insert(&self.data_ks, &item_k, enc_val);
let next_kind = if target_expire == 0 {
HashFieldStateKind::Persistent
} else {
HashFieldStateKind::LiveTTL
};
state_cache.insert(
f_bytes,
CachedFieldState {
kind: next_kind,
expire: target_expire,
raw: None,
},
);
}
if meta.base.size == 0 {
if metadata_existed {
batch.remove(&self.meta_ks, meta_k.as_bytes());
}
} else if meta_changed || !metadata_existed {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
Ok(true)
}
pub fn hsetex<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_vals: &[(F, V)],
flags: &[HSetEx],
) -> Result<bool> {
let now_ms = current_now_ms();
let opts = HashSetExOptions::from_flags(flags, now_ms);
self.set_fields_with_expire(key, field_vals, opts)
}
pub fn get_fields_with_expire<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
options: HashGetExOptions,
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(vec![None; fields.len()]),
};
if meta.is_legacy_subkey_encoding() {
return Err(Error::invalid_data(
ERR_HASH_FIELD_EXPIRATION_LEGACY_ENCODING,
));
}
let mut results = Vec::with_capacity(fields.len());
let mut batch = self.db.batch();
let mut meta_changed = false;
let mut state_cache: HashMap<&[u8], CachedFieldState> =
HashMap::with_capacity_and_hasher(fields.len(), Default::default());
let is_immediate = options.ttl_action == TTLAction::Set
&& is_immediate_expire(options.expire_at_ms, now_ms);
for f in fields {
let f_bytes = f.as_ref();
let item_k = kc.hash_key_bytes(k_str, f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let raw_opt = self.data_ks.get(&item_k)?;
let state_entry = match raw_opt {
Some(raw) => {
if let Some(state) = decode_field_state(&meta, &raw, now_ms) {
CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
}
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
}
}
None => CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
};
state_cache.insert(f_bytes, state_entry.clone());
state_entry
};
match entry.kind {
HashFieldStateKind::Missing => {
results.push(None);
}
HashFieldStateKind::ExpiredTTLPhysical => {
results.push(None);
batch.remove(&self.data_ks, &item_k);
meta.apply_ttl_to_deleted();
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
}
HashFieldStateKind::Persistent | HashFieldStateKind::LiveTTL => {
let payload = entry
.raw
.as_ref()
.and_then(|s| meta.decode_subkey_value(s))
.map(|(_, p)| p.to_vec())
.unwrap_or_default();
results.push(Some(payload.clone()));
match options.ttl_action {
TTLAction::Discard | TTLAction::Keep => {}
TTLAction::Persist => {
if entry.kind == HashFieldStateKind::LiveTTL {
meta.apply_ttl_to_persistent();
meta_changed = true;
let enc = meta.encode_subkey_value(&payload, 0);
batch.insert(&self.data_ks, &item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Persistent,
expire: 0,
raw: entry.raw,
},
);
}
}
TTLAction::Set => {
if is_immediate {
batch.remove(&self.data_ks, &item_k);
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_deleted();
} else {
meta.apply_ttl_to_deleted();
}
meta_changed = true;
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
},
);
} else {
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_ttl(options.expire_at_ms);
} else if entry.expire != options.expire_at_ms {
meta.apply_ttl_to_ttl(options.expire_at_ms);
}
meta_changed = true;
let enc = meta.encode_subkey_value(&payload, options.expire_at_ms);
batch.insert(&self.data_ks, &item_k, enc);
state_cache.insert(
f_bytes,
CachedFieldState {
kind: HashFieldStateKind::LiveTTL,
expire: options.expire_at_ms,
raw: entry.raw,
},
);
}
}
}
}
}
}
if meta_changed {
if meta.base.size == 0 {
batch.remove(&self.meta_ks, meta_k.as_bytes());
} else {
batch.insert(&self.meta_ks, meta_k.as_bytes(), meta.encode());
}
batch.commit()?;
}
Ok(results)
}
pub fn hgetex<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
flag: Option<HGetEx>,
) -> Result<Option<Vec<u8>>> {
let now_ms = current_now_ms();
let opts = match flag {
Some(f) => HashGetExOptions::from_flag(f, now_ms),
None => HashGetExOptions::default(),
};
let mut vals = self.get_fields_with_expire(key, &[field], opts)?;
Ok(vals.pop().unwrap_or(None))
}
pub fn hmgetex<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
fields: &[F],
flag: Option<HGetEx>,
) -> Result<Vec<Option<Vec<u8>>>> {
let now_ms = current_now_ms();
let opts = match flag {
Some(f) => HashGetExOptions::from_flag(f, now_ms),
None => HashGetExOptions::default(),
};
self.get_fields_with_expire(key, fields, opts)
}
pub fn hrandfield<K: AsRef<[u8]>>(
&self,
key: K,
count: i64,
with_values: bool,
) -> Result<Vec<HashRandField>> {
if count == 0 {
return Ok(Vec::new());
}
if with_values {
let all_pairs = self.hgetall(&key)?;
if all_pairs.is_empty() {
return Ok(Vec::new());
}
if count > 0 {
let sample_cnt = (count as usize).min(all_pairs.len());
if sample_cnt == 1 {
let idx = fastrand::usize(0..all_pairs.len());
let (ref f, ref v) = all_pairs[idx];
return Ok(vec![(f.clone(), Some(v.clone()))]);
}
let mut indices: Vec<usize> = (0..all_pairs.len()).collect();
for i in 0..sample_cnt {
let j = fastrand::usize(i..all_pairs.len());
indices.swap(i, j);
}
let mut out = Vec::with_capacity(sample_cnt);
for &idx in &indices[..sample_cnt] {
let (ref f, ref v) = all_pairs[idx];
out.push((f.clone(), Some(v.clone())));
}
Ok(out)
} else {
let total = count.unsigned_abs() as usize;
let mut out = Vec::with_capacity(total);
for _ in 0..total {
let idx = fastrand::usize(0..all_pairs.len());
let (ref f, ref v) = all_pairs[idx];
out.push((f.clone(), Some(v.clone())));
}
Ok(out)
}
} else {
let all_keys = self.hkeys(&key)?;
if all_keys.is_empty() {
return Ok(Vec::new());
}
if count > 0 {
let sample_cnt = (count as usize).min(all_keys.len());
if sample_cnt == 1 {
let idx = fastrand::usize(0..all_keys.len());
return Ok(vec![(all_keys[idx].clone(), None)]);
}
let mut indices: Vec<usize> = (0..all_keys.len()).collect();
for i in 0..sample_cnt {
let j = fastrand::usize(i..all_keys.len());
indices.swap(i, j);
}
let mut out = Vec::with_capacity(sample_cnt);
for &idx in &indices[..sample_cnt] {
out.push((all_keys[idx].clone(), None));
}
Ok(out)
} else {
let total = count.unsigned_abs() as usize;
let mut out = Vec::with_capacity(total);
for _ in 0..total {
let idx = fastrand::usize(0..all_keys.len());
out.push((all_keys[idx].clone(), None));
}
Ok(out)
}
}
}
pub fn hrangebylex<K: AsRef<[u8]>>(
&self,
key: K,
spec: RangeLexSpec,
) -> Result<Vec<HashFieldPair>> {
let kc = KeyComposer::new("default");
let k_str = str::from_utf8(key.as_ref()).unwrap_or("");
let meta_k = kc.hash_meta(k_str);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta(&meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
let prefix = kc.hash_prefix(k_str);
let mut matching = Vec::new();
for g in self.data_ks.prefix(&prefix) {
let (k, v) = g.into_inner()?;
if !k.starts_with(&prefix) {
break;
}
let field_bytes = &k[prefix.len()..];
let min_ok = if spec.min_infinite {
true
} else if spec.minex {
field_bytes > spec.min.as_slice()
} else {
field_bytes >= spec.min.as_slice()
};
if !min_ok {
continue;
}
let max_ok = if spec.max_infinite {
true
} else if spec.maxex {
field_bytes < spec.max.as_slice()
} else {
field_bytes <= spec.max.as_slice()
};
if !max_ok {
break;
}
if let Some((exp, payload)) = meta.decode_subkey_value(&v)
&& !is_field_expired(exp, now_ms)
{
matching.push((field_bytes.to_vec(), payload.to_vec()));
}
}
if spec.reversed {
matching.reverse();
}
let start = spec.offset.min(matching.len());
let limit = spec.count.unwrap_or(matching.len());
let end = (start + limit).min(matching.len());
Ok(matching[start..end].to_vec())
}
pub fn hscan<K: AsRef<[u8]>>(
&self,
key: K,
cursor: usize,
limit: usize,
pattern: Option<&[u8]>,
) -> Result<HashScanResult> {
let all = self.hgetall(key)?;
let is_match_all = match pattern {
Some(p) => p == b"*",
None => true,
};
let matched: Vec<(Vec<u8>, Vec<u8>)> = if is_match_all {
all
} else {
let pat = pattern.unwrap_or(b"*");
all.into_iter()
.filter(|(f, _)| matches_glob_bytes(pat, f))
.collect()
};
let start = cursor.min(matched.len());
let end = (start + limit).min(matched.len());
let next_cursor = if end >= matched.len() { 0 } else { end };
let slice = if start < matched.len() {
matched[start..end].to_vec()
} else {
Vec::new()
};
Ok((next_cursor, slice))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_wedb_hash_basic_ops() -> Result<()> {
let dir = tempfile::tempdir().map_err(Error::Io)?;
let db = WeDb::open(dir.path())?;
let added = db.hset("myhash", &[("f1", "v1"), ("f2", "v2"), ("f3", "v3")])?;
assert_eq!(added, 3);
assert_eq!(db.hlen("myhash")?, 3);
assert_eq!(db.hlen_with_mode("myhash", HashLengthMode::Approximate)?, 3);
let v1 = db.hget("myhash", "f1")?;
assert_eq!(v1, Some(b"v1".to_vec()));
assert!(db.hexists("myhash", "f1")?);
assert!(!db.hexists("myhash", "f99")?);
assert_eq!(db.hstrlen("myhash", "f1")?, 2);
let dup_added = db.hset("myhash", &[("f1", "v1_new"), ("f1", "v1_final")])?;
assert_eq!(dup_added, 0);
assert_eq!(db.hget("myhash", "f1")?, Some(b"v1_final".to_vec()));
assert!(!db.hsetnx("myhash", "f1", "v_nx")?);
assert!(db.hsetnx("myhash", "f4", "v4")?);
assert_eq!(db.hlen("myhash")?, 4);
let mvals = db.hmget("myhash", &["f1", "f2", "f_none"])?;
assert_eq!(
mvals,
vec![Some(b"v1_final".to_vec()), Some(b"v2".to_vec()), None]
);
let all = db.hgetall("myhash")?;
assert_eq!(all.len(), 4);
let keys = db.hkeys("myhash")?;
assert_eq!(keys.len(), 4);
let vals = db.hvals("myhash")?;
assert_eq!(vals.len(), 4);
db.hset("h_num", &[("c1", "10"), ("f1", "2.5")])?;
let inc = db.hincrby("h_num", "c1", 5)?;
assert_eq!(inc, 15);
let inc_f = db.hincrbyfloat("h_num", "f1", 1.5)?;
assert!((inc_f - 4.0).abs() < 1e-6);
let deleted = db.hdel("myhash", &["f1", "f2", "f1"])?;
assert_eq!(deleted, 2);
assert_eq!(db.hlen("myhash")?, 2);
Ok(())
}
#[test]
fn test_wedb_hash_expiration_ttl() -> Result<()> {
let dir = tempfile::tempdir().map_err(Error::Io)?;
let db = WeDb::open(dir.path())?;
db.hset("htest", &[("a", "1"), ("b", "2"), ("c", "3")])?;
let res = db.hexpire("htest", &["a", "b"], 10, HExpire::None)?;
assert_eq!(res, vec![1, 1]);
let ttls = db.httl("htest", &["a", "b", "c", "d"])?;
assert!(ttls[0] > 0 && ttls[0] <= 10);
assert!(ttls[1] > 0 && ttls[1] <= 10);
assert_eq!(ttls[2], HASH_FIELD_PERSISTENT);
assert_eq!(ttls[3], HASH_FIELD_NOT_FOUND);
let pttls = db.hpttl("htest", &["a", "c"])?;
assert!(pttls[0] > 0);
assert_eq!(pttls[1], HASH_FIELD_PERSISTENT);
let ex_times = db.hexpiretime("htest", &["a", "c"])?;
assert!(ex_times[0] > 0);
assert_eq!(ex_times[1], HASH_FIELD_PERSISTENT);
let cond_nx = db.hexpire("htest", &["a", "c"], 20, HExpire::Nx)?;
assert_eq!(cond_nx, vec![0, 1]);
let persist_res = db.hpersist("htest", &["a", "b"])?;
assert_eq!(persist_res, vec![1, 1]);
assert_eq!(
db.httl("htest", &["a", "b"])?,
vec![HASH_FIELD_PERSISTENT, HASH_FIELD_PERSISTENT]
);
let imm = db.hexpire("htest", &["a"], 0, HExpire::None)?;
assert_eq!(imm, vec![HASH_EXPIRE_DELETED]);
assert!(!db.hexists("htest", "a")?);
assert_eq!(db.hlen("htest")?, 2);
Ok(())
}
#[test]
fn test_wedb_hash_hgetdel_and_setex_getex() -> Result<()> {
let dir = tempfile::tempdir().map_err(Error::Io)?;
let db = WeDb::open(dir.path())?;
let opt_nx = HashSetExOptions::new(
HashFieldSetCondition::Fnx,
TTLAction::Set,
current_now_ms() + 5000,
);
let ok = db.set_fields_with_expire("hex_test", &[("k1", "v1"), ("k2", "v2")], opt_nx)?;
assert!(ok);
assert_eq!(db.hlen("hex_test")?, 2);
let fail_nx = db.set_fields_with_expire("hex_test", &[("k1", "v1_new")], opt_nx)?;
assert!(!fail_nx);
let ok2 = db.hsetex("hex_test2", &[("a", "100")], &[HSetEx::Ex(10), HSetEx::Fnx])?;
assert!(ok2);
assert_eq!(db.hlen("hex_test2")?, 1);
let get_opt = HashGetExOptions::persist();
let gvals = db.get_fields_with_expire("hex_test", &["k1", "k_none"], get_opt)?;
assert_eq!(gvals, vec![Some(b"v1".to_vec()), None]);
assert_eq!(db.httl("hex_test", &["k1"])?, vec![HASH_FIELD_PERSISTENT]);
let gval_single = db.hgetex("hex_test2", "a", Some(HGetEx::Persist))?;
assert_eq!(gval_single, Some(b"100".to_vec()));
assert_eq!(db.httl("hex_test2", &["a"])?, vec![HASH_FIELD_PERSISTENT]);
let del_vals = db.hgetdel("hex_test", &["k1", "k2", "k_none"])?;
assert_eq!(
del_vals,
vec![Some(b"v1".to_vec()), Some(b"v2".to_vec()), None]
);
assert_eq!(db.hlen("hex_test")?, 0);
Ok(())
}
#[test]
fn test_wedb_hash_randfield_lex_scan() -> Result<()> {
let dir = tempfile::tempdir().map_err(Error::Io)?;
let db = WeDb::open(dir.path())?;
db.hset(
"h_lex",
&[("f1", "v1"), ("f2", "v2"), ("f3", "v3"), ("f4", "v4")],
)?;
let rand_single = db.hrandfield("h_lex", 1, false)?;
assert_eq!(rand_single.len(), 1);
let rand_with_vals = db.hrandfield("h_lex", 2, true)?;
assert_eq!(rand_with_vals.len(), 2);
assert!(rand_with_vals[0].1.is_some());
let lex_spec = RangeLexSpec {
min: b"f1".to_vec(),
max: b"f3".to_vec(),
minex: false,
maxex: false,
min_infinite: false,
max_infinite: false,
offset: 0,
count: None,
reversed: false,
};
let lex_res = db.hrangebylex("h_lex", lex_spec)?;
assert_eq!(lex_res.len(), 3);
assert_eq!(lex_res[0].0, b"f1");
assert_eq!(lex_res[2].0, b"f3");
let (cursor, scanned) = db.hscan("h_lex", 0, 10, Some(b"f*"))?;
assert_eq!(cursor, 0);
assert_eq!(scanned.len(), 4);
Ok(())
}
}