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,
ERR_WRONG_TYPE, 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, HashItemKeyComposer, HashMeta,
HashSubkeyEncodingMode, compose_hash_meta_key, compose_hash_prefix, compose_hash_prefix_stack,
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::{HashMapExt, HashSetExt, RapidHashMap as HashMap, RapidHashSet as HashSet};
use crate::db::WeDb;
use crate::error::{Error, Result};
use crate::key_composer::{KeyComposer, KeyTag, matches_glob_bytes};
use crate::meta::{current_now_ms, parse_redis_float, parse_redis_integer};
use crate::string::format_float_bytes;
#[inline(always)]
const fn ceil_div_1000(val: u64) -> u64 {
val.div_ceil(1000)
}
#[inline]
fn parse_hash_integer(v: &[u8]) -> Result<i64> {
parse_redis_integer(v, ERR_HASH_VALUE_NOT_INTEGER)
}
#[inline]
fn parse_hash_float(v: &[u8]) -> Result<f64> {
parse_redis_float(v, ERR_HASH_VALUE_NOT_FLOAT)
}
#[derive(Clone)]
struct CachedFieldState {
kind: HashFieldStateKind,
expire: u64,
raw: Option<fjall::Slice>,
}
impl WeDb {
#[inline]
fn get_hash_meta_checked(
&self,
kc: &KeyComposer<'_>,
k: &[u8],
mk: &[u8],
now: u64,
) -> Result<Option<HashMeta>> {
self.get_meta_checked(kc, k, mk, now)
}
#[inline]
fn load_field_state(
&self,
meta: &HashMeta,
item_k: &[u8],
now_ms: u64,
) -> Result<CachedFieldState> {
match self.data.get(item_k)? {
Some(raw) => {
if let Some(state) = decode_field_state(meta, &raw, now_ms) {
Ok(CachedFieldState {
kind: state.kind,
expire: state.expire,
raw: Some(raw),
})
} else {
Ok(CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
})
}
}
None => Ok(CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}),
}
}
#[inline]
fn prepare_hash_meta_for_write(
&self,
kc: &KeyComposer<'_>,
k_bytes: &[u8],
meta_k: &[u8],
now_ms: u64,
batch: &mut fjall::OwnedWriteBatch,
) -> Result<(HashMeta, bool)> {
match self.meta.get(meta_k)? {
Some(m_bytes) => match HashMeta::decode(&m_bytes) {
Some(meta) if !meta.is_expired(now_ms) => Ok((meta, true)),
Some(_) => {
let prefix = compose_hash_prefix_stack(kc, k_bytes);
self.clear_prefix_in_batch(&prefix, batch)?;
Ok((HashMeta::new_with_version(0, 0), false))
}
None => Ok((HashMeta::new_with_version(0, 0), false)),
},
None => {
self.check_key_not_other_type(kc, k_bytes, KeyTag::HashMeta.as_slice(), now_ms)?;
Ok((HashMeta::new_with_version(0, 0), false))
}
}
}
#[inline]
pub fn hiter<K: AsRef<[u8]>, F>(&self, key: K, mut f: F) -> Result<()>
where
F: FnMut(&[u8], &[u8]) -> bool,
{
let key_bytes = key.as_ref();
let kc = KeyComposer::new("default");
let meta_k = compose_hash_meta_key(&kc, key_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(&kc, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(()),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(());
}
let prefix = compose_hash_prefix(&kc, key_bytes);
for g in self.data.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)
{
let field = &k[prefix.len()..];
if !f(field, payload) {
break;
}
}
}
Ok(())
}
pub fn hset<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
key: K,
field_vals: &[(F, V)],
) -> Result<usize> {
self.hset_with_kc(&KeyComposer::new("default"), key, field_vals)
}
pub fn hset_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field_vals: &[(F, V)],
) -> Result<usize> {
if field_vals.is_empty() {
return Ok(0);
}
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
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_bytes, &meta_k, now_ms, &mut batch)?;
let mut added = 0;
let mut meta_changed = false;
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
if field_vals.len() == 1 {
let f_bytes = field_vals[0].0.as_ref();
let v_bytes = field_vals[0].1.as_ref();
let item_k = composer.key_for_field(f_bytes);
let kind = if metadata_existed
&& let Some(raw) = self.data.get(item_k)?
&& let Some(state) = decode_field_state(&meta, &raw, now_ms)
{
if state.kind == HashFieldStateKind::Persistent && state.value == v_bytes {
return Ok(0);
}
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, item_k, enc_val);
} else {
let mut seen = HashSet::with_capacity(field_vals.len());
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();
for (f_bytes, v_bytes) in unique_items {
let item_k = composer.key_for_field(f_bytes);
let kind = if metadata_existed
&& let Some(raw) = self.data.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, item_k, enc_val);
}
}
if meta_changed || !metadata_existed {
batch.insert(&self.meta, &meta_k, 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_with_kc(&KeyComposer::new("default"), key, field_vals)?;
Ok(())
}
#[inline]
pub fn hmset_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>, V: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field_vals: &[(F, V)],
) -> Result<()> {
self.hset_with_kc(kc, 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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
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_bytes, &meta_k, now_ms, &mut batch)?;
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
let f_bytes = field.as_ref();
let item_k = composer.key_for_field(f_bytes);
let kind = if metadata_existed
&& let Some(raw) = self.data.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, item_k, enc_val);
batch.insert(&self.meta, &meta_k, meta.encode());
batch.commit()?;
Ok(true)
}
pub fn hget<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
key: K,
field: F,
) -> Result<Option<Vec<u8>>> {
self.hget_with_kc(&KeyComposer::new("default"), key, field)
}
pub fn hget_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field: F,
) -> Result<Option<Vec<u8>>> {
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, k_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(None),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(None);
}
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
let item_k = composer.key_for_field(field.as_ref());
match self.data.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> {
self.hdel_with_kc(&KeyComposer::new("default"), key, fields)
}
pub fn hdel_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
fields: &[F],
) -> Result<usize> {
if fields.is_empty() {
return Ok(0);
}
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta_checked(kc, k_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(0),
};
let mut seen = HashSet::with_capacity(fields.len());
let mut deleted = 0;
let mut physical_removed = 0;
let mut batch = self.db.batch();
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
for f in fields {
let f_bytes = f.as_ref();
if !seen.insert(f_bytes) {
continue;
}
let item_k = composer.key_for_field(f_bytes);
if let Some(raw) = self.data.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, item_k);
}
HashFieldStateKind::LiveTTL => {
deleted += 1;
physical_removed += 1;
meta.apply_ttl_to_deleted();
batch.remove(&self.data, item_k);
}
HashFieldStateKind::ExpiredTTLPhysical => {
physical_removed += 1;
meta.apply_ttl_to_deleted();
batch.remove(&self.data, item_k);
}
HashFieldStateKind::Missing => {}
}
}
}
if physical_removed == 0 {
return Ok(0);
}
if meta.base.size == 0 {
batch.remove(&self.meta, &meta_k);
} else {
batch.insert(&self.meta, &meta_k, meta.encode());
}
batch.commit()?;
Ok(deleted)
}
#[inline]
pub fn hexists<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<bool> {
self.hexists_with_kc(&KeyComposer::new("default"), key, field)
}
pub fn hexists_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field: F,
) -> Result<bool> {
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, k_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(false),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(false);
}
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
let item_k = composer.key_for_field(field.as_ref());
if let Some(raw) = self.data.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_kc<K: AsRef<[u8]>>(&self, kc: &KeyComposer<'_>, key: K) -> Result<usize> {
self.hlen_with_mode_and_kc(kc, key, HashLengthMode::Accurate)
}
pub fn hlen_with_mode<K: AsRef<[u8]>>(&self, key: K, mode: HashLengthMode) -> Result<usize> {
self.hlen_with_mode_and_kc(&KeyComposer::new("default"), key, mode)
}
pub fn hlen_with_mode_and_kc<K: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
mode: HashLengthMode,
) -> Result<usize> {
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta_checked(kc, k_bytes, &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_bytes, &meta_k, &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, &meta_k);
batch.commit()?;
return Ok(0);
}
self.scan_and_repair_hash(kc, k_bytes, &meta_k, &mut meta, now_ms)
}
fn scan_and_repair_hash(
&self,
kc: &KeyComposer<'_>,
k_bytes: &[u8],
meta_k: &[u8],
meta: &mut HashMeta,
now_ms: u64,
) -> Result<usize> {
let prefix = compose_hash_prefix(kc, k_bytes);
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.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, &*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, meta_k);
} else {
repaired.clear_bounds_if_no_ttl_candidates();
batch.insert(&self.meta, meta_k, 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>>>> {
self.hmget_with_kc(&KeyComposer::new("default"), key, fields)
}
pub fn hmget_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
fields: &[F],
) -> Result<Vec<Option<Vec<u8>>>> {
if fields.is_empty() {
return Ok(Vec::new());
}
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, k_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(vec![None; fields.len()]),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(vec![None; fields.len()]);
}
let mut results = Vec::with_capacity(fields.len());
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
for f in fields {
let item_k = composer.key_for_field(f.as_ref());
let val = match self.data.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>)>> {
match fetch_type {
HashFetchType::All => self.hgetall(key),
HashFetchType::OnlyKey => {
let keys = self.hkeys(key)?;
Ok(keys.into_iter().map(|k| (k, Vec::new())).collect())
}
HashFetchType::OnlyValue => {
let vals = self.hvals(key)?;
Ok(vals.into_iter().map(|v| (Vec::new(), v)).collect())
}
}
}
pub fn hgetall<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
self.hgetall_with_kc(&KeyComposer::new("default"), key)
}
pub fn hgetall_with_kc<K: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
) -> Result<Vec<(Vec<u8>, Vec<u8>)>> {
let key_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(Vec::new()),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(Vec::new());
}
let prefix = compose_hash_prefix(kc, key_bytes);
let mut items = Vec::with_capacity(meta.base.size as usize);
for g in self.data.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)
{
let field_bytes = &k[prefix.len()..];
items.push((field_bytes.to_vec(), payload.to_vec()));
}
}
Ok(items)
}
pub fn hkeys<K: AsRef<[u8]>>(&self, key: K) -> Result<Vec<Vec<u8>>> {
self.hkeys_with_kc(&KeyComposer::new("default"), key)
}
pub fn hkeys_with_kc<K: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
) -> Result<Vec<Vec<u8>>> {
let key_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(Vec::new()),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(Vec::new());
}
let prefix = compose_hash_prefix(kc, key_bytes);
let mut keys = Vec::with_capacity(meta.base.size as usize);
for g in self.data.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>>> {
self.hvals_with_kc(&KeyComposer::new("default"), key)
}
pub fn hvals_with_kc<K: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
) -> Result<Vec<Vec<u8>>> {
let key_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, key_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, key_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(Vec::new()),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(Vec::new());
}
let prefix = compose_hash_prefix(kc, key_bytes);
let mut vals = Vec::with_capacity(meta.base.size as usize);
for g in self.data.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> {
self.hincrby_with_kc(&KeyComposer::new("default"), key, field, increment)
}
pub fn hincrby_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field: F,
increment: i64,
) -> Result<i64> {
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
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_bytes, &meta_k, now_ms, &mut batch)?;
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
let item_k = composer.key_for_field(field.as_ref());
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.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 {
cur_num = parse_hash_integer(payload)?;
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, item_k, enc_val);
if meta_changed || !metadata_existed {
batch.insert(&self.meta, &meta_k, 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> {
self.hincrbyfloat_with_kc(&KeyComposer::new("default"), key, field, increment)
}
pub fn hincrbyfloat_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field: F,
increment: f64,
) -> Result<f64> {
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
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_bytes, &meta_k, now_ms, &mut batch)?;
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
let item_k = composer.key_for_field(field.as_ref());
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.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 {
cur_num = parse_hash_float(payload)?;
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 = zmij::Buffer::new();
let num_bytes = format_float_bytes(new_num, &mut buf);
let enc_val = meta.encode_subkey_value(num_bytes, keep_expire);
batch.insert(&self.data, item_k, enc_val);
if meta_changed || !metadata_existed {
batch.insert(&self.meta, &meta_k, meta.encode());
}
batch.commit()?;
Ok(new_num)
}
#[inline]
pub fn hstrlen<K: AsRef<[u8]>, F: AsRef<[u8]>>(&self, key: K, field: F) -> Result<usize> {
self.hstrlen_with_kc(&KeyComposer::new("default"), key, field)
}
pub fn hstrlen_with_kc<K: AsRef<[u8]>, F: AsRef<[u8]>>(
&self,
kc: &KeyComposer<'_>,
key: K,
field: F,
) -> Result<usize> {
let k_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(kc, k_bytes, &meta_k, now_ms)? {
Some(m) if m.base.size > 0 => m,
_ => return Ok(0),
};
if meta.upper != 0 && now_ms > meta.upper && meta.persist == 0 {
return Ok(0);
}
let mut composer = HashItemKeyComposer::new(kc, k_bytes);
let item_k = composer.key_for_field(field.as_ref());
if let Some(raw) = self.data.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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let mut meta = match self.get_hash_meta_checked(&kc, k_bytes, &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(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = self.load_field_state(&meta, item_k, now_ms)?;
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, 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, 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, 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, &meta_k);
} else {
batch.insert(&self.meta, &meta_k, 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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(&kc, k_bytes, &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 result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
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 = composer.key_for_field(f_bytes);
let res = match self.data.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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(&kc, k_bytes, &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 result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
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 = composer.key_for_field(f_bytes);
let res = match self.data.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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(&kc, k_bytes, &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 result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
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 = composer.key_for_field(f_bytes);
let res = match self.data.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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(&kc, k_bytes, &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 result_cache: HashMap<&[u8], i64> = HashMap::with_capacity(fields.len());
let mut results = Vec::with_capacity(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
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 = composer.key_for_field(f_bytes);
let res = match self.data.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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta_checked(&kc, k_bytes, &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(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = self.load_field_state(&meta, item_k, now_ms)?;
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, 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, 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, &meta_k);
} else {
batch.insert(&self.meta, &meta_k, 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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta_checked(&kc, k_bytes, &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(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
for f in fields {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = self.load_field_state(&meta, item_k, now_ms)?;
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, 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, 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, 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, &meta_k);
} else {
batch.insert(&self.meta, &meta_k, 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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
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_bytes, &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(field_values.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
if options.condition != HashFieldSetCondition::None {
for (f, _) in field_values {
let f_bytes = f.as_ref();
let item_k = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else if metadata_existed {
let state_entry = self.load_field_state(&meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
};
if entry.kind == HashFieldStateKind::ExpiredTTLPhysical {
batch.remove(&self.data, 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, &meta_k);
} else {
batch.insert(&self.meta, &meta_k, meta.encode());
}
batch.commit()?;
}
return Ok(false);
}
}
}
let mut seen = HashSet::with_capacity(field_values.len());
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 = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else if metadata_existed {
let state_entry = self.load_field_state(&meta, item_k, now_ms)?;
state_cache.insert(f_bytes, state_entry.clone());
state_entry
} else {
CachedFieldState {
kind: HashFieldStateKind::Missing,
expire: 0,
raw: None,
}
};
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, 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, 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, &meta_k);
}
} else if meta_changed || !metadata_existed {
batch.insert(&self.meta, &meta_k, 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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let mut meta = match self.get_hash_meta_checked(&kc, k_bytes, &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(fields.len());
let mut composer = HashItemKeyComposer::new(&kc, k_bytes);
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 = composer.key_for_field(f_bytes);
let entry = if let Some(cached) = state_cache.get(f_bytes) {
cached.clone()
} else {
let state_entry = self.load_field_state(&meta, item_k, now_ms)?;
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, 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, 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, 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::LiveTTL
&& entry.expire == options.expire_at_ms
{
continue;
}
if entry.kind == HashFieldStateKind::Persistent {
meta.apply_persistent_to_ttl(options.expire_at_ms);
} else {
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, 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, &meta_k);
} else {
batch.insert(&self.meta, &meta_k, 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 mut all_pairs = self.hgetall(&key)?;
if all_pairs.is_empty() {
return Ok(Vec::new());
}
let total = all_pairs.len();
if count > 0 {
let sample_cnt = (count as usize).min(total);
if sample_cnt == 1 {
let idx = fastrand::usize(0..total);
let (f, v) = all_pairs.swap_remove(idx);
return Ok(vec![(f, Some(v))]);
}
if sample_cnt == total {
fastrand::shuffle(&mut all_pairs);
return Ok(all_pairs.into_iter().map(|(f, v)| (f, Some(v))).collect());
}
for i in 0..sample_cnt {
let j = fastrand::usize(i..total);
all_pairs.swap(i, j);
}
all_pairs.truncate(sample_cnt);
Ok(all_pairs.into_iter().map(|(f, v)| (f, Some(v))).collect())
} else {
let total_sample = count.unsigned_abs() as usize;
let mut out = Vec::with_capacity(total_sample);
for _ in 0..total_sample {
let idx = fastrand::usize(0..total);
let (ref f, ref v) = all_pairs[idx];
out.push((f.clone(), Some(v.clone())));
}
Ok(out)
}
} else {
let mut all_keys = self.hkeys(&key)?;
if all_keys.is_empty() {
return Ok(Vec::new());
}
let total = all_keys.len();
if count > 0 {
let sample_cnt = (count as usize).min(total);
if sample_cnt == 1 {
let idx = fastrand::usize(0..total);
let f = all_keys.swap_remove(idx);
return Ok(vec![(f, None)]);
}
if sample_cnt == total {
fastrand::shuffle(&mut all_keys);
return Ok(all_keys.into_iter().map(|f| (f, None)).collect());
}
for i in 0..sample_cnt {
let j = fastrand::usize(i..total);
all_keys.swap(i, j);
}
all_keys.truncate(sample_cnt);
Ok(all_keys.into_iter().map(|f| (f, None)).collect())
} else {
let total_sample = count.unsigned_abs() as usize;
let mut out = Vec::with_capacity(total_sample);
for _ in 0..total_sample {
let idx = fastrand::usize(0..total);
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_bytes = key.as_ref();
let meta_k = compose_hash_meta_key(&kc, k_bytes);
let now_ms = current_now_ms();
let meta = match self.get_hash_meta_checked(&kc, k_bytes, &meta_k, now_ms)? {
Some(m) => m,
None => return Ok(Vec::new()),
};
if spec.count == Some(0) {
return Ok(Vec::new());
}
let prefix = compose_hash_prefix(&kc, k_bytes);
if !spec.reversed {
let limit = spec.count.unwrap_or(usize::MAX);
let mut matching = Vec::with_capacity(limit.min(128));
let mut skipped = 0;
for g in self.data.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)
{
if skipped < spec.offset {
skipped += 1;
continue;
}
matching.push((field_bytes.to_vec(), payload.to_vec()));
if matching.len() >= limit {
break;
}
}
}
return Ok(matching);
}
let mut matching = Vec::new();
for g in self.data.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()));
}
}
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 is_match_all = match pattern {
Some(p) => p == b"*",
None => true,
};
let pat = pattern.unwrap_or(b"*");
let mut skipped = 0;
let mut matched = Vec::with_capacity(limit);
let mut has_more = false;
self.hiter(key, |field, value| {
if is_match_all || matches_glob_bytes(pat, field) {
if skipped < cursor {
skipped += 1;
} else if matched.len() < limit {
matched.push((field.to_vec(), value.to_vec()));
} else {
has_more = true;
return false;
}
}
true
})?;
let next_cursor = if has_more { cursor + matched.len() } else { 0 };
Ok((next_cursor, matched))
}
}